{"record":{"id":"0d107a8c307edb75","repo":"risingwavelabs/risingwave","slug":"primary-key-column-not-found-in-iceberg-schema","errorCode":null,"errorMessage":"Primary key column {} not found in Iceberg schema","messagePattern":"Primary key column (.+?) not found in Iceberg schema","errorType":"validation","errorClass":"SinkError::Config","httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/iceberg/writer.rs","lineNumber":316,"sourceCode":"    ) -> Self {\n        Self::Created(IcebergSinkWriterArgs {\n            config,\n            sink_param,\n            writer_param,\n            upsert_primary_key_column_names,\n        })\n    }\n}\n\npub(super) fn resolve_equality_delete_field_ids(\n    primary_key_column_names: &[String],\n    iceberg_schema: &iceberg::spec::Schema,\n) -> Result<Vec<i32>> {\n    primary_key_column_names\n        .iter()\n        .map(|column_name| {\n            iceberg_schema.field_id_by_name(column_name).ok_or_else(|| {\n                SinkError::Config(anyhow!(\n                    \"Primary key column {} not found in Iceberg schema\",\n                    column_name\n                ))\n            })\n        })\n        .collect()\n}\n\nimpl IcebergSinkWriterInner {\n    pub fn build_append_only(\n        config: &IcebergConfig,\n        table: Table,\n        writer_param: &SinkWriterParam,\n    ) -> Result<Self> {\n        let SinkWriterParam {\n            extra_partition_col_idx,\n            actor_id,\n            sink_id,","sourceCodeStart":298,"sourceCodeEnd":334,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/iceberg/writer.rs#L298-L334","documentation":"Thrown by resolve_equality_delete_field_ids when a primary key column name from the sink's PRIMARY KEY definition cannot be resolved to a field id in the target Iceberg table schema. Iceberg identifies columns by integer field ids, not names, so an unresolvable name means the sink configuration does not match the actual table schema and equality-delete files cannot be projected.","triggerScenarios":"Calling build_upsert for an Iceberg sink whose primary_key_column_names include a column missing (or renamed/case-mismatched) in table.metadata().current_schema(); sinking into a table evolved to drop a PK column.","commonSituations":"Typos or case differences between PRIMARY KEY in the CREATE SINK statement and the Iceberg table columns; pointing the sink at a different table than planned; Iceberg schema evolution removing a key column after sink creation.","solutions":["Compare each primary key column name against the Iceberg table schema (iceberg_table_metadata.current_schema()) and fix the sink's PRIMARY KEY to use exact existing column names","Check for case sensitivity: Iceberg field lookup is exact, so normalize quoting/casing in the sink definition","Re-create the sink against the correct table, or re-evolve the table to restore the dropped PK column","Add a pre-flight validation in the sink creation path that resolves all PK columns to field ids before the writer is built"],"exampleFix":"// before\niceberg_schema.field_id_by_name(column_name).ok_or_else(|| {\n    SinkError::Config(anyhow!(\"Primary key column {} not found in Iceberg schema\", column_name))\n})\n// after\niceberg_schema.field_id_by_name(column_name).ok_or_else(|| {\n    let known = iceberg_schema.fields().iter().map(|f| f.name()).collect::<Vec<_>>().join(\", \");\n    SinkError::Config(anyhow!(\n        \"Primary key column {} not found in Iceberg schema (available: [{}])\",\n        column_name,\n        known\n    ))\n})","handlingStrategy":"validation","validationCode":"for name in primary_key_column_names {\n    if iceberg_schema.field_id_by_name(name).is_none() {\n        return Err(format!(\"PK column {} missing from Iceberg schema\", name));\n    }\n}","typeGuard":"fn pk_columns_resolve(names: &[String], schema: &iceberg::spec::Schema) -> bool {\n    names.iter().all(|n| schema.field_id_by_name(n).is_some())\n}","tryCatchPattern":"let ids = resolve_equality_delete_field_ids(&pk_names, &schema)\n    .map_err(|e| format!(\"sink PK does not match table schema: {e}\"))?;","preventionTips":["Validate PK column names against the Iceberg schema at sink creation time","Watch for case sensitivity when quoting column names in SQL","Re-check the schema after any Iceberg table evolution (ALTER/DROP COLUMN)","Keep the sink definition and the target table in the same migration/change review"],"tags":["rust","iceberg","schema","primary-key"],"backgroundTag":"resource-not-found","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"}