{"record":{"id":"2bf5d756ce45bf0e","repo":"risingwavelabs/risingwave","slug":"iceberg-pk-index-writer-expected-checkpoint","errorCode":null,"errorMessage":"iceberg pk-index writer {} expected checkpoint {:?} starting at {}, got {:?}","messagePattern":"iceberg pk-index writer (.+?) expected checkpoint (.+?) starting at (.+?), got (.+?)","errorType":"exception","errorClass":"StreamExecutorError","httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/iceberg_with_pk_index/writer.rs","lineNumber":316,"sourceCode":"                    barrier.epoch,\n                    self.sink_id,\n                    self.ctx.id,\n                    PbIcebergPkIndexSinkRole::Writer,\n                    Some(metadata),\n                );\n        }\n        yield Message::Barrier(barrier);\n    }\n\n    fn validate_compaction_barrier(\n        &self,\n        barrier: &Barrier,\n        expected_task: IcebergCompactionTaskId,\n        expected_phase: Phase,\n        expected_prev: u64,\n    ) -> StreamExecutorResult<()> {\n        if !barrier.is_checkpoint() || barrier.epoch.prev != expected_prev {\n            bail!(\n                \"iceberg pk-index writer {} expected checkpoint {:?} starting at {}, got {:?}\",\n                self.sink_id,\n                expected_phase,\n                expected_prev,\n                barrier\n            );\n        }\n        match barrier.iceberg_pk_index_compaction() {\n            Some(context)\n                if context.sink_id == self.sink_id\n                    && context.task_id == expected_task\n                    && context.phase == expected_phase as i32 =>\n            {\n                Ok(())\n            }\n            _ => bail!(\n                \"iceberg pk-index writer {} expected matching {:?} context for task {}, got {:?}\",\n                self.sink_id,","sourceCodeStart":298,"sourceCodeEnd":334,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/iceberg_with_pk_index/writer.rs#L298-L334","documentation":"During compaction phases the writer validates that each incoming barrier is a checkpoint whose epoch.prev equals the epoch the writer expects to start from for the given phase. This error means the barrier either wasn't a checkpoint or its epoch chain didn't start at expected_prev, so the writer's phase machine cannot proceed safely.","triggerScenarios":"Raised in validate_compaction_barrier (called from execute_resolving_right) when !barrier.is_checkpoint() || barrier.epoch.prev != expected_prev — a non-checkpoint barrier, or a checkpoint whose prev epoch diverges from the tracked expected_prev, arrives while awaiting a compaction phase transition.","commonSituations":"Epoch divergence after failover/recovery (barriers replayed from a different epoch chain); a checkpoint skipped or reordered; meta issuing compaction phase transitions inconsistent with the writer's local epoch tracking; scaled-in/actor migration changing barrier flow.","solutions":["Compare the logged expected_prev with the barrier's epoch.prev to find where the chain diverged.","Restart the fragment from the last consistent checkpoint so epochs realign.","Check meta's compaction scheduler for phase transitions issued on wrong barriers.","Verify no configuration change (checkpoint interval, actor layout) occurred mid-compaction.","File a bug with both expected and observed barriers if reproducible."],"exampleFix":null,"handlingStrategy":"validation","validationCode":"// check barrier properties before invoking the phase machine\nassert!(barrier.is_checkpoint(), \"phase transitions require checkpoints\");\nassert_eq!(barrier.epoch.prev, expected_prev, \"epoch chain diverged\");","typeGuard":"fn barrier_matches(barrier: &Barrier, expected_prev: u64) -> bool {\n    barrier.is_checkpoint() && barrier.epoch.prev == expected_prev\n}","tryCatchPattern":"if let Err(e) = writer.execute_resolving_right(msg).await {\n    if e.to_string().contains(\"expected checkpoint\") {\n        error!(\"epoch divergence during compaction; recovering from last checkpoint\");\n    }\n    return Err(e);\n}","preventionTips":["Recover fragments from checkpoints rather than replaying raw barriers.","Avoid changing actor layout or checkpoint interval mid-compaction.","Monitor epoch chain continuity between meta and stream nodes."],"tags":["iceberg","compaction","barrier","checkpoint","epoch"],"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"}