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
- Align meta and all worker binaries to the same version so the position-delete commit-result protobuf schema matches.
- Inspect the raw report metadata for the failing sink to verify the payload type is IcebergPositionDeleteCommitResult.
- Restart the sink/fragment so workers rebuild fresh reports instead of replaying stale metadata.
- 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
- Keep meta and worker protobuf definitions from the same generated source.
- Verify the report role matches the payload type before serializing.
- Avoid replaying persisted report bytes across version upgrades.
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
- DecodeError {error}
- iceberg pk-index sink report missing metadata in aggregate_r
- iceberg pk-index sink report has invalid role: {}
- Pb decode error: {0}
- Failed to decode prost: field not found `{}`
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/117447bd302a6a92.
Report an issue: GitHub.