{"record":{"id":"c4c3ff3e4ba8461d","repo":"risingwavelabs/risingwave","slug":"iceberg-pk-index-sink-reports-disagree-on-schema-i","errorCode":null,"errorMessage":"iceberg pk-index sink reports disagree on schema_id: {} vs {}","messagePattern":"iceberg pk-index sink reports disagree on schema_id: (.+?) vs (.+?)","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/meta/src/manager/iceberg_pk_index_sink/coordinator.rs","lineNumber":628,"sourceCode":"\n    Ok(IcebergPkIndexSinkAggResult {\n        schema_id: shared_schema_id.unwrap(),\n        partition_spec_id: shared_partition_spec_id.unwrap(),\n        data_files,\n        delete_files,\n        overwrite_files,\n    })\n}\n\nfn align_report_id(\n    schema_id: i32,\n    partition_spec_id: i32,\n    shared_schema_id: &mut Option<i32>,\n    shared_partition_spec_id: &mut Option<i32>,\n) -> Result<()> {\n    match shared_schema_id {\n        Some(prev) if *prev != schema_id => {\n            bail!(\n                \"iceberg pk-index sink reports disagree on schema_id: {} vs {}\",\n                prev,\n                schema_id\n            );\n        }\n        None => *shared_schema_id = Some(schema_id),\n        _ => {}\n    }\n    match shared_partition_spec_id {\n        Some(prev) if *prev != partition_spec_id => {\n            bail!(\n                \"iceberg pk-index sink reports disagree on partition_spec_id: {} vs {}\",\n                prev,\n                partition_spec_id\n            );\n        }\n        None => *shared_partition_spec_id = Some(partition_spec_id),\n        _ => {}","sourceCodeStart":610,"sourceCodeEnd":646,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/manager/iceberg_pk_index_sink/coordinator.rs#L610-L646","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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."],"exampleFix":null,"handlingStrategy":"validation","validationCode":"// before triggering commit, ensure all workers report the same schema id\nlet ids: HashSet<i32> = reports.iter().map(|r| r.schema_id).collect();\nif ids.len() > 1 {\n    return Err(anyhow!(\"workers disagree on schema_id: {:?}\", ids));\n}","typeGuard":null,"tryCatchPattern":"// retry once after workers refresh schema\nif let Err(e) = commit().await {\n    if e.to_string().contains(\"disagree on schema_id\") {\n        refresh_worker_schemas().await;\n        commit().await?;\n    } else {\n        return Err(e);\n    }\n}","preventionTips":["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."],"tags":["iceberg","schema","consistency"],"backgroundTag":"schema-validation-failed","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}