risingwavelabs/risingwave · error
iceberg pk-index sink reports disagree on partition_spec_id
Error message
iceberg pk-index sink reports disagree on partition_spec_id: {} vs {} What it means
Analogous to the schema check: every report in a single commit must agree on the table's partition_spec_id. align_report_id compares each report's id against the shared value and bails with both ids when they differ, preventing a commit that mixes files partitioned under different specs.
Solutions
- Identify the worker with the stale partition_spec_id from logs and restart it to reload current table metadata.
- Force a refresh of Iceberg table metadata on all workers so partition specs align before the next commit.
- Confirm the intended spec id on the Iceberg table and roll back an accidental spec change if unintended.
- Restart the sink fragment so all reports regenerate against the same spec.
Defensive patterns
Strategy: validation
Validate before calling
// before triggering commit, ensure all workers report the same partition spec id
let ids: HashSet<i32> = reports.iter().map(|r| r.partition_spec_id).collect();
if ids.len() > 1 {
return Err(anyhow!("workers disagree on partition_spec_id: {:?}", ids));
} Try / catch
// retry once after refreshing table metadata
if let Err(e) = commit().await {
if e.to_string().contains("disagree on partition_spec_id") {
refresh_table_metadata().await;
commit().await?;
} else {
return Err(e);
}
} Prevention
- Freeze partition-spec changes while a sink commit is in flight.
- Refresh Iceberg table metadata on workers promptly after spec evolution.
- Avoid restoring workers from checkpoints older than the latest spec change.
When it happens
Trigger: pre_commit_epoch -> aggregate_reports -> align_report_id when two worker reports for the same sink/epoch carry different partition_spec_id values, typically around a partition-spec evolution.
Common situations: Partition spec evolution applied while workers still hold the old spec; a worker restored from an old checkpoint/metadata snapshot; replayed stale reports from before the spec change.
Understand the failure class
Background: Schema validation failed / invalid input schema: payload rejected because its shape doesn't match the expected schema — this error's family across 28 libraries.
Related errors
- iceberg pk-index sink reports disagree on schema_id
- Invalid partition fields
- schema_id and partition_spec_id should be the same in all…
- The `n` must be set with `bucket` and `truncate`
- adlsgen2.authority_host does not parse as a URL
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/c636446dfab905eb.
Report an issue: GitHub.
Appendix: source
Thrown at src/meta/src/manager/iceberg_pk_index_sink/coordinator.rs:639
schema_id: i32,
partition_spec_id: i32,
shared_schema_id: &mut Option<i32>,
shared_partition_spec_id: &mut Option<i32>,
) -> Result<()> {
match shared_schema_id {
Some(prev) if *prev != schema_id => {
bail!(
"iceberg pk-index sink reports disagree on schema_id: {} vs {}",
prev,
schema_id
);
}
None => *shared_schema_id = Some(schema_id),
_ => {}
}
match shared_partition_spec_id {
Some(prev) if *prev != partition_spec_id => {
bail!(
"iceberg pk-index sink reports disagree on partition_spec_id: {} vs {}",
prev,
partition_spec_id
);
}
None => *shared_partition_spec_id = Some(partition_spec_id),
_ => {}
}
Ok(())
}
fn encode_pre_commit_state(
agg_result: &IcebergPkIndexSinkAggResult,
snapshot_id: i64,
) -> Result<Vec<u8>> {
let agg_result = serde_json::to_vec(agg_result)?;
Ok(PbIcebergPkIndexPreCommitState {
agg_result,View on GitHub (pinned to 6469eb736d)