{"record":{"id":"c636446dfab905eb","repo":"risingwavelabs/risingwave","slug":"iceberg-pk-index-sink-reports-disagree-on-partitio","errorCode":null,"errorMessage":"iceberg pk-index sink reports disagree on partition_spec_id: {} vs {}","messagePattern":"iceberg pk-index sink reports disagree on partition_spec_id: (.+?) vs (.+?)","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/meta/src/manager/iceberg_pk_index_sink/coordinator.rs","lineNumber":639,"sourceCode":"    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        _ => {}\n    }\n    Ok(())\n}\n\nfn encode_pre_commit_state(\n    agg_result: &IcebergPkIndexSinkAggResult,\n    snapshot_id: i64,\n) -> Result<Vec<u8>> {\n    let agg_result = serde_json::to_vec(agg_result)?;\n    Ok(PbIcebergPkIndexPreCommitState {\n        agg_result,","sourceCodeStart":621,"sourceCodeEnd":657,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/manager/iceberg_pk_index_sink/coordinator.rs#L621-L657","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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."],"exampleFix":null,"handlingStrategy":"validation","validationCode":"// before triggering commit, ensure all workers report the same partition spec id\nlet ids: HashSet<i32> = reports.iter().map(|r| r.partition_spec_id).collect();\nif ids.len() > 1 {\n    return Err(anyhow!(\"workers disagree on partition_spec_id: {:?}\", ids));\n}","typeGuard":null,"tryCatchPattern":"// retry once after refreshing table metadata\nif let Err(e) = commit().await {\n    if e.to_string().contains(\"disagree on partition_spec_id\") {\n        refresh_table_metadata().await;\n        commit().await?;\n    } else {\n        return Err(e);\n    }\n}","preventionTips":["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."],"tags":["iceberg","partitioning","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"}