{"record":{"id":"d13ade53a5e1294d","repo":"risingwavelabs/risingwave","slug":"compaction-resolver-pk-column-missing-column-desc","errorCode":null,"errorMessage":"compaction resolver PK column missing column_desc","messagePattern":"compaction resolver PK column missing column_desc","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/from_proto/iceberg_with_pk_index/compaction_resolver.rs","lineNumber":69,"sourceCode":"        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);\n        let local_barrier_manager = params.local_barrier_manager.clone();\n        let meta_client = params.env.meta_client().ok_or_else(|| {\n            anyhow!(\"meta client is required for iceberg pk-index compaction resolver\")\n        })?;\n        let exec = CompactionResolverExecutor::new(\n            params.actor_context,\n            sink_id,\n            iceberg_config,\n            pk_indices,\n            pk_data_types,\n            params.config.developer.chunk_size,\n            local_barrier_manager,","sourceCodeStart":51,"sourceCodeEnd":87,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/from_proto/iceberg_with_pk_index/compaction_resolver.rs#L51-L87","documentation":"When building PK data types for the Iceberg compaction resolver, each entry in `node.pk_columns` must carry a `column_desc` from which the `ColumnDesc` (and thus data type) is derived. A missing `column_desc` makes the PK type unknowable, so `new_boxed_executor` errors with this message.","triggerScenarios":"A proto node whose `pk_columns[i].column_desc` is `None` during the `.map(...ok_or_else(...))` collection, i.e. partially serialized PK column metadata reaching the executor conversion.","commonSituations":"Proto serialization bugs dropping nested column_desc; schema evolution/version skew where old plans lack column descriptors; corrupted fragment graph during recovery or manual replay.","solutions":["Regenerate the stream plan so each PK column includes its column_desc.","Upgrade frontend/meta together with the stream node to fix proto serialization mismatches.","Re-create the sink so the fragment carries complete column descriptors."],"exampleFix":null,"handlingStrategy":"validation","validationCode":"// before dispatching the plan\nassert!(\n    node.pk_columns.iter().all(|c| c.column_desc.is_some()),\n    \"every PK column must carry a column_desc\"\n);","typeGuard":"fn pk_descs_present(node: &StreamNode) -> bool {\n    node.pk_columns.iter().all(|c| c.column_desc.is_some())\n}","tryCatchPattern":null,"preventionTips":["Verify proto serialization round-trips preserve nested column_desc","Add schema-completeness checks in the fragment graph builder","Upgrade all components together to avoid plan-schema skew"],"tags":["streaming","iceberg","schema","protobuf"],"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-14T16:17:12.679Z"}