risingwavelabs/risingwave · error
iceberg pk-index sink reports disagree on schema_id
Error message
iceberg pk-index sink reports disagree on schema_id: {} vs {} What it means
All worker reports in one Iceberg pk-index sink commit must reference the same table schema. align_report_id records the first observed schema_id and bails if a later report disagrees, because committing files written under different schemas would corrupt the Iceberg table. The error reports both conflicting ids.
Solutions
- Identify the lagging worker from logs and restart it so it reloads the current table schema.
- Ensure schema-change notifications fully propagate; wait for the next epoch so all workers align on the new schema.
- Check the Iceberg table's schema history to confirm the expected schema_id and reconcile workers pinned to an outdated snapshot.
- If a worker is stuck with a cached schema, restart the sink fragment.
Defensive patterns
Strategy: validation
Validate before calling
// before triggering commit, ensure all workers report the same schema id
let ids: HashSet<i32> = reports.iter().map(|r| r.schema_id).collect();
if ids.len() > 1 {
return Err(anyhow!("workers disagree on schema_id: {:?}", ids));
} Try / catch
// retry once after workers refresh schema
if let Err(e) = commit().await {
if e.to_string().contains("disagree on schema_id") {
refresh_worker_schemas().await;
commit().await?;
} else {
return Err(e);
}
} Prevention
- Hold commits during schema evolution until all workers acknowledge the new schema.
- Monitor worker schema-version drift in metrics.
- Restart workers that miss schema-change notifications.
When it happens
Trigger: pre_commit_epoch -> aggregate_reports -> align_report_id when two reports for the same sink/epoch carry different schema_id values, e.g. one worker still on the pre-evolution schema and another on the post-evolution one.
Common situations: Schema evolution landing mid-epoch while some writers/mergers still run with the old cached schema; stale workers that missed the schema-change notification; replayed old reports merged with fresh ones.
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
- Can't find schema by id
- Column order mismatch at position
- compaction resolver PK column missing column_desc
- error computing partition type
- Field not found in our schema
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/c4c3ff3e4ba8461d.
Report an issue: GitHub.
Appendix: source
Thrown at src/meta/src/manager/iceberg_pk_index_sink/coordinator.rs:628
Ok(IcebergPkIndexSinkAggResult {
schema_id: shared_schema_id.unwrap(),
partition_spec_id: shared_partition_spec_id.unwrap(),
data_files,
delete_files,
overwrite_files,
})
}
fn align_report_id(
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),
_ => {}View on GitHub (pinned to 6469eb736d)