{"record":{"id":"44079962ce2fdb5e","repo":"risingwavelabs/risingwave","slug":"iceberg-pk-index-writer-expected-resolver-inser","errorCode":null,"errorMessage":"iceberg pk-index writer {} expected resolver inserts, got {op:?}","messagePattern":"iceberg pk-index writer (.+?) expected resolver inserts, got (.+?)","errorType":"exception","errorClass":"StreamExecutorError","httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/iceberg_with_pk_index/writer.rs","lineNumber":264,"sourceCode":"            for (pk, pos) in insert_pks.into_iter().zip_eq_fast(positions) {\n                let mut index_row_data = Vec::with_capacity(pk_indices.len() + 2);\n                for datum in pk.iter() {\n                    index_row_data.push(datum);\n                }\n                index_row_data.push(Some(ScalarRefImpl::Utf8(&pos.path)));\n                index_row_data.push(Some(ScalarRefImpl::Int64(pos.pos)));\n                self.pk_index_state_table.insert(index_row_data.as_slice());\n            }\n        }\n\n        self.delete_position_buffer = Some(delete_position_buffer);\n        self.pk_index_state_table.try_flush().await?;\n    }\n\n    async fn apply_resolver_chunk(&mut self, chunk: StreamChunk) -> StreamExecutorResult<()> {\n        for (op, row) in chunk.rows() {\n            if op != Op::Insert {\n                bail!(\n                    \"iceberg pk-index writer {} expected resolver inserts, got {op:?}\",\n                    self.sink_id\n                );\n            }\n            self.pk_index_state_table.insert(row);\n        }\n        self.pk_index_state_table.try_flush().await?;\n        Ok(())\n    }\n\n    #[try_stream(ok = Message, error = StreamExecutorError)]\n    async fn checkpoint_barrier(&mut self, barrier: Barrier) {\n        barrier.assume_no_update_vnode_bitmap(self.ctx.id)?;\n        let mut metadata = None;\n        if barrier.is_checkpoint() {\n            if let Some(chunk) = self\n                .delete_position_buffer\n                .take()","sourceCodeStart":246,"sourceCodeEnd":282,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/iceberg_with_pk_index/writer.rs#L246-L282","documentation":"The resolver side of the pk-index writer only expects Insert operations: the equi-join resolver produces the rows that must be inserted into the pk-index state table. Any Update/Delete op reaching apply_resolver_chunk violates the protocol and aborts with this error. The state table is append-only from this path, so non-inserts are unrecoverable here.","triggerScenarios":"Raised in apply_resolver_chunk (called from execute_resolving_right) when a StreamChunk row from the resolving (right) side carries Op::UpdateDelete, Op::UpdateInsert, or Op::Delete instead of Op::Insert.","commonSituations":"Upstream executor changes emitting CDC-style updates into the resolver side; incorrect join/executor wiring that passes mutating ops to the writer's resolver input; data replay path differences after recovery.","solutions":["Inspect the upstream executor feeding the resolver side to see why non-insert ops are emitted.","Verify executor wiring/version compatibility for the iceberg pk-index sink fragment.","If replay/recovery produced this, restart the fragment from a clean checkpoint.","File a bug including the op variant and actor id if it occurs with standard pipelines."],"exampleFix":"// before\nif op != Op::Insert {\n    bail!(\"... expected resolver inserts, got {op:?}\");\n}\n// after\n// keep as-is: fix the upstream instead — e.g. ensure the materialize/resolver input\n// is configured with INSERT-only output (no update/delete propagation).","handlingStrategy":"validation","validationCode":"// filter or assert at the producer side feeding the resolver input\nfor (op, row) in chunk.rows() {\n    assert_eq!(op, Op::Insert, \"resolver input must be insert-only\");\n}","typeGuard":"fn is_insert_only(chunk: &StreamChunk) -> bool {\n    chunk.rows().all(|(op, _)| op == Op::Insert)\n}","tryCatchPattern":"if let Err(e) = writer.execute_resolving_right(chunk).await {\n    if e.to_string().contains(\"expected resolver inserts\") {\n        error!(\"non-insert op reached resolver input; check upstream executor wiring\");\n    }\n    return Err(e);\n}","preventionTips":["Verify the resolver-side upstream is configured for insert-only output.","Add a debug assertion on op type in executor tests.","Keep executor topology changes covered by fragment integration tests."],"tags":["iceberg","stream-chunk","op-type","protocol-violation","streaming-executor"],"backgroundTag":"invalid-enum-value","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"}