{"record":{"id":"4c0e4b7de60fb87d","repo":"risingwavelabs/risingwave","slug":"materialize-pk-index-sink-serializeddatafile","errorCode":null,"errorMessage":"materialize pk-index sink SerializedDataFile","messagePattern":"materialize pk-index sink SerializedDataFile","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/meta/src/manager/iceberg_pk_index_sink/coordinator.rs","lineNumber":500,"sourceCode":"                    .snapshots()\n                    .any(|s| s.snapshot_id() == snapshot_id)\n                {\n                    return Ok(table);\n                }\n\n                let schema = table.metadata().current_schema();\n                let partition_type =\n                    resolve_partition_type(&table, merged.partition_spec_id, schema)\n                        .map_err(|e| CommitError::Commit(anyhow!(e)))?;\n\n                let materialize =\n                    |serialized: &SerializedDataFile| -> Result<DataFile, CommitError> {\n                        serialized\n                            .clone()\n                            .try_into(merged.partition_spec_id, &partition_type, schema)\n                            .map_err(|err| {\n                                CommitError::Commit(\n                                    anyhow!(err)\n                                        .context(\"materialize pk-index sink SerializedDataFile\"),\n                                )\n                            })\n                    };\n\n                // Reuse the files already materialized during pre-commit backfill when available;\n                // otherwise (recovery / unpartitioned) materialize from the persisted form once.\n                let add_files: Vec<DataFile> = match materialized_add_files {\n                    Some(add_files) => add_files,\n                    None => merged\n                        .data_files\n                        .iter()\n                        .chain(merged.delete_files.iter())\n                        .map(&materialize)\n                        .collect::<Result<Vec<_>, _>>()?,\n                };\n                let overwrite_files: Vec<DataFile> = merged\n                    .overwrite_files","sourceCodeStart":482,"sourceCodeEnd":518,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/manager/iceberg_pk_index_sink/coordinator.rs#L482-L518","documentation":"During commit, serialized data files persisted in the aggregation result must be materialized back into Iceberg `DataFile`s against the table's current schema and partition type. `SerializedDataFile::try_into` validates partition values and schema fields; a mismatch throws this error, wrapped as a `CommitError::Commit`.","triggerScenarios":"In `commit_one_epoch`, the persisted `SerializedDataFile` cannot convert with the current `partition_spec_id` / `partition_type` / schema — e.g. table schema evolved or partition spec changed between pre-commit persistence and commit, or the blob was written against a different table version.","commonSituations":"External schema evolution (ALTER TABLE on the Iceberg table by another engine) while a commit is pending; recovered pending rows from an older table snapshot; partition spec replaced in the catalog.","solutions":["Inspect the underlying conversion error in the log to see which field/partition value failed validation.","Check whether the Iceberg table's schema or partition spec changed externally after the report was produced; align/revert the external change or recompute the commit.","Abort the stale pending epoch so it can be re-run with reports matching the current table schema."],"exampleFix":"// avoid external evolution of a sink-managed table:\n// before: ALTER TABLE db.tbl SET PARTITION SPEC (...); -- by external engine\n// after: coordinate spec changes with RisingWave sink downtime","handlingStrategy":"try-catch","validationCode":"// verify schema/spec stability before commit\nanyhow::ensure!(\n    table.metadata().current_schema().schema_id() == merged.schema_id,\n    \"table schema changed since pre-commit\"\n);\nanyhow::ensure!(\n    table.metadata().default_partition_spec_id() == merged.partition_spec_id,\n    \"partition spec changed since pre-commit\"\n);","typeGuard":null,"tryCatchPattern":"match serialized.clone().try_into(spec_id, &partition_type, schema) {\n    Ok(df) => { /* use df */ }\n    Err(e) => return Err(CommitError::Commit(anyhow!(e).context(\"materialize pk-index sink SerializedDataFile\"))),\n}","preventionTips":["Avoid external ALTER TABLE / partition-spec changes on sink-managed tables.","Keep pending epochs short so schema skew windows stay small.","Add a pre-commit check that table schema_id/partition_spec_id match the reports."],"tags":["iceberg","schema-mismatch","commit","type-mismatch"],"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"}