risingwavelabs/risingwave · error

decode pk-index sink position-delete merger metadata

Error message

decode pk-index sink position-delete merger metadata

What it means

A PositionDeleteMerger worker's commit report carries opaque metadata bytes that the coordinator decodes into IcebergPositionDeleteCommitResult via try_from(meta). When those bytes cannot be decoded (wrong message type, truncated, or an incompatible schema/version), aggregate_reports wraps the failure with this context and fails the pre-commit phase.

Source

Thrown at src/meta/src/manager/iceberg_pk_index_sink/coordinator.rs:596

            .ok()
            .filter(|r| !matches!(r, PbIcebergPkIndexSinkRole::Unspecified))
            .ok_or_else(|| anyhow!("iceberg pk-index sink report has invalid role: {}", r.role))?;

        match role {
            PbIcebergPkIndexSinkRole::Writer => {
                let commit_result = IcebergCommitResult::try_from(meta)?;
                align_report_id(
                    commit_result.schema_id,
                    commit_result.partition_spec_id,
                    &mut shared_schema_id,
                    &mut shared_partition_spec_id,
                )?;
                data_files.extend(commit_result.data_files);
            }
            PbIcebergPkIndexSinkRole::PositionDeleteMerger => {
                let commit_result =
                    IcebergPositionDeleteCommitResult::try_from(meta).map_err(|e| {
                        anyhow!(e).context("decode pk-index sink position-delete merger metadata")
                    })?;
                align_report_id(
                    commit_result.schema_id,
                    commit_result.partition_spec_id,
                    &mut shared_schema_id,
                    &mut shared_partition_spec_id,
                )?;
                delete_files.extend(commit_result.delete_files);
                overwrite_files.extend(commit_result.overwrite_files);
            }
            _ => unreachable!(),
        }
    }

    Ok(IcebergPkIndexSinkAggResult {
        schema_id: shared_schema_id.unwrap(),
        partition_spec_id: shared_partition_spec_id.unwrap(),
        data_files,

View on GitHub (pinned to 6469eb736d)

Solutions

  1. Align meta and all worker binaries to the same version so the position-delete commit-result protobuf schema matches.
  2. Inspect the raw report metadata for the failing sink to verify the payload type is IcebergPositionDeleteCommitResult.
  3. Restart the sink/fragment so workers rebuild fresh reports instead of replaying stale metadata.
  4. Check the producing worker to confirm it serializes the position-delete result, not the writer result.

Example fix

// before
let commit_result = IcebergPositionDeleteCommitResult::try_from(meta)?;
// after (fail with diagnosis)
let commit_result = IcebergPositionDeleteCommitResult::try_from(meta.clone())
    .context("report meta is not IcebergPositionDeleteCommitResult; check worker/meta version alignment")?;
Defensive patterns

Strategy: validation

Validate before calling

// before trusting report meta
if report.meta.is_none() || report.meta.as_ref().unwrap().is_empty() {
    return Err(anyhow!("position-delete merger report meta is empty"));
}

Type guard

fn decodable_position_delete_meta(meta: &[u8]) -> bool {
    IcebergPositionDeleteCommitResult::try_from(meta.to_vec()).is_ok()
}

Try / catch

// caller
match coordinator.pre_commit_epoch(epoch).await {
    Err(e) if e.to_string().contains("position-delete merger metadata") => {
        // version skew / bad payload: log and re-register the sink
    }
    other => other?,
}

Prevention

When it happens

Trigger: pre_commit_epoch -> aggregate_reports handles a report with role PositionDeleteMerger and calls IcebergPositionDeleteCommitResult::try_from(meta), which errors because the report's meta bytes are not a valid position-delete commit-result payload.

Common situations: Version skew between meta and workers producing different protobuf layouts; a Writer-role payload delivered under the PositionDeleteMerger role; corrupted or replayed report metadata after a meta restart.

Understand the failure class

Background: "cannot parse invalid wire-format data", "cannot unmarshal", "failed unmarshalling": protobuf unmarshal errors explained — this error's family across 10 libraries.

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/117447bd302a6a92. Report an issue: GitHub.