{"record":{"id":"1e20de0b7a399c0d","repo":"risingwavelabs/risingwave","slug":"missing-primary-key-columns-in-compaction-resolver","errorCode":null,"errorMessage":"missing primary-key columns in compaction resolver","messagePattern":"missing primary-key columns in compaction resolver","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/from_proto/iceberg_with_pk_index/compaction_resolver.rs","lineNumber":57,"sourceCode":"        assert!(\n            params.input.is_empty(),\n            \"compaction resolver executor should not have input\"\n        );\n\n        let sink_id = node.sink_id;\n\n        let properties_with_secret = LocalSecretManager::global()\n            .fill_secrets(node.properties.clone(), node.secret_refs.clone())?;\n        let iceberg_config = IcebergConfig::from_btreemap(properties_with_secret)\n            .map_err(|err| StreamExecutorError::from((err, sink_id)))?;\n\n        let pk_indices = node\n            .pk_columns\n            .iter()\n            .map(|column| column.data_file_index as usize)\n            .collect::<Vec<_>>();\n        if pk_indices.is_empty() {\n            return Err(anyhow!(\"missing primary-key columns in compaction resolver\").into());\n        }\n\n        let pk_data_types = node\n            .pk_columns\n            .iter()\n            .map(|column| {\n                column\n                    .column_desc\n                    .as_ref()\n                    .map(ColumnDesc::from)\n                    .map(|column| column.data_type)\n                    .ok_or_else(|| anyhow!(\"compaction resolver PK column missing column_desc\"))\n            })\n            .collect::<Result<Vec<_>, _>>()?;\n\n        let barrier_receiver = params\n            .local_barrier_manager\n            .subscribe_barrier(params.actor_context.id);","sourceCodeStart":39,"sourceCodeEnd":75,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/from_proto/iceberg_with_pk_index/compaction_resolver.rs#L39-L75","documentation":"The Iceberg compaction resolver conversion derives primary key indices from `node.pk_columns` and requires at least one entry. An empty list means the resolver cannot identify row identity for resolving compaction results, so `new_boxed_executor` fails early with this message.","triggerScenarios":"Building the compaction resolver executor when `node.pk_columns` is empty (or all entries map to zero data-file indices yielding an empty vec), i.e. the proto node was produced without PK columns for a with-pk-index Iceberg sink.","commonSituations":"Planning an Iceberg sink with PK index but omitting primary key definition; frontend serialization bug dropping pk_columns; older plan versions lacking the pk_columns field.","solutions":["Define primary key columns on the Iceberg sink/table so pk_columns are populated in the proto.","Upgrade frontend/meta so pk_columns are serialized for iceberg_with_pk_index nodes.","Recreate the sink materialization to regenerate the fragment with PK columns."],"exampleFix":"-- before\nCREATE SINK s FROM mv; -- no primary key\n-- after\nCREATE SINK s FROM mv WITH PRIMARY KEY (id); -- pk_columns now non-empty","handlingStrategy":"validation","validationCode":"// before building the node proto\nassert!(!pk_columns.is_empty(), \"iceberg with-pk-index sink requires primary key columns\");","typeGuard":"fn has_pk(node: &StreamNode) -> bool { !node.pk_columns.is_empty() }","tryCatchPattern":null,"preventionTips":["Always define a primary key for iceberg_with_pk_index sinks","Add plan validation that pk_columns is non-empty at planning time","Check proto serialization after frontend changes"],"tags":["streaming","iceberg","primary-key"],"backgroundTag":"missing-required-argument","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"}