{"record":{"id":"ea842cd8113a521a","repo":"risingwavelabs/risingwave","slug":"iceberg-pk-index-writer-expected-begin-context","errorCode":null,"errorMessage":"iceberg pk-index writer {} expected Begin context for task {}, got {:?}","messagePattern":"iceberg pk-index writer (.+?) expected Begin context for task (.+?), got (.+?)","errorType":"exception","errorClass":"StreamExecutorError","httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/iceberg_with_pk_index/writer.rs","lineNumber":374,"sourceCode":"    }\n\n    fn compaction_begin(\n        &self,\n        barrier: &Barrier,\n    ) -> StreamExecutorResult<Option<IcebergCompactionTaskId>> {\n        let context = match barrier.iceberg_pk_index_compaction() {\n            Some(context) if context.sink_id == self.sink_id => context,\n            _ => return Ok(None),\n        };\n        if context.phase == Phase::End as i32 {\n            bail!(\n                \"iceberg pk-index writer {} received unexpected End in Normal mode for task {}\",\n                self.sink_id,\n                context.task_id\n            );\n        }\n        if context.phase != Phase::Begin as i32 {\n            bail!(\n                \"iceberg pk-index writer {} expected Begin context for task {}, got {:?}\",\n                self.sink_id,\n                context.task_id,\n                context.phase\n            );\n        }\n        Ok(Some(context.task_id))\n    }\n\n    #[try_stream(ok = Message, error = StreamExecutorError)]\n    async fn execute_inner(mut self) {\n        let mut input = self.input.take().unwrap().execute();\n        let mut resolver_input = self.resolver_input.take().unwrap().execute();\n\n        // Consume the first barrier.\n        let barrier = expect_first_barrier(&mut input).await?;\n        let remap_first = expect_first_barrier(&mut resolver_input).await?;\n        self.validate_aligned_barriers(&barrier, &remap_first)?;","sourceCodeStart":356,"sourceCodeEnd":392,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/iceberg_with_pk_index/writer.rs#L356-L392","documentation":"When compaction_begin receives a barrier carrying a compaction context for this sink in Normal mode, that context must be Phase::Begin (entering the compaction flow). Any other phase (Apply/End) indicates a lifecycle inconsistency: the writer is being asked to start a phase that can only occur mid-compaction.","triggerScenarios":"Raised in compaction_begin (called from execute_normal) when context.phase != Phase::Begin as i32 — e.g. a Phase::Apply or Phase::End context barrier arrives while the writer is in Normal mode for the task.","commonSituations":"Missed or mis-ordered Begin barrier after recovery; meta issuing Apply/End phases to a writer that never entered the task; racing phase transitions when tasks are rescheduled across actors; node version skew.","solutions":["Identify the observed phase (in the error) vs the expected Begin and trace which barrier was missed.","Restart the fragment from the last checkpoint to reset the compaction state machine.","Audit meta's compaction phase scheduling for this task_id to ensure Begin precedes Apply/End.","Check for actor migration or scale changes that dropped the Begin barrier.","File a bug with the task id and phase if reproducible."],"exampleFix":null,"handlingStrategy":"validation","validationCode":"// only dispatch Begin-phase contexts into compaction_begin in Normal mode\nif ctx.phase == Phase::Begin as i32 {\n    writer.compaction_begin(barrier, task).await?;\n}","typeGuard":"fn is_begin_phase(ctx: &IcebergPkIndexCompactionContext) -> bool {\n    ctx.phase == Phase::Begin as i32\n}","tryCatchPattern":"match res {\n    Err(e) if e.to_string().contains(\"expected Begin context\") => {\n        error!(\"compaction phases out of order; restarting fragment from last checkpoint\");\n    }\n    other => other?,\n}","preventionTips":["Ensure meta emits compaction phases strictly in Begin → Apply → End order.","Reset writer compaction state on recovery to Normal mode before applying new Begins.","Audit task rescheduling so phases never leak to a writer that missed Begin."],"tags":["iceberg","compaction","barrier","phase","state-machine"],"backgroundTag":"invalid-state-transition","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"}