risingwavelabs/risingwave · error · StreamExecutorError
iceberg pk-index writer {} expected checkpoint {:?} starting
Error message
iceberg pk-index writer {} expected checkpoint {:?} starting at {}, got {:?} What it means
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.
Source
Thrown at src/stream/src/executor/iceberg_with_pk_index/writer.rs:316
barrier.epoch,
self.sink_id,
self.ctx.id,
PbIcebergPkIndexSinkRole::Writer,
Some(metadata),
);
}
yield Message::Barrier(barrier);
}
fn validate_compaction_barrier(
&self,
barrier: &Barrier,
expected_task: IcebergCompactionTaskId,
expected_phase: Phase,
expected_prev: u64,
) -> StreamExecutorResult<()> {
if !barrier.is_checkpoint() || barrier.epoch.prev != expected_prev {
bail!(
"iceberg pk-index writer {} expected checkpoint {:?} starting at {}, got {:?}",
self.sink_id,
expected_phase,
expected_prev,
barrier
);
}
match barrier.iceberg_pk_index_compaction() {
Some(context)
if context.sink_id == self.sink_id
&& context.task_id == expected_task
&& context.phase == expected_phase as i32 =>
{
Ok(())
}
_ => bail!(
"iceberg pk-index writer {} expected matching {:?} context for task {}, got {:?}",
self.sink_id,View on GitHub (pinned to 6469eb736d)
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.
Defensive patterns
Strategy: validation
Validate before calling
// check barrier properties before invoking the phase machine assert!(barrier.is_checkpoint(), "phase transitions require checkpoints"); assert_eq!(barrier.epoch.prev, expected_prev, "epoch chain diverged");
Type guard
fn barrier_matches(barrier: &Barrier, expected_prev: u64) -> bool {
barrier.is_checkpoint() && barrier.epoch.prev == expected_prev
} Try / catch
if let Err(e) = writer.execute_resolving_right(msg).await {
if e.to_string().contains("expected checkpoint") {
error!("epoch divergence during compaction; recovering from last checkpoint");
}
return Err(e);
} Prevention
- 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.
When it happens
Trigger: 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.
Common situations: 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.
Understand the failure class
Background: "Invalid state transition" errors: "status must be X, actually Y", "already rejected/charging/uninstalled", "cannot ... while running" — what they mean when a library rejects your call — this error's family across 31 libraries.
Related errors
- iceberg pk-index merger {} received a non-checkpoint compact
- iceberg pk-index writer {} expected matching {:?} context fo
- iceberg pk-index writer {} received unexpected End in Normal
- iceberg pk-index writer {} expected Begin context for task {
- Iceberg source should not have input executor!
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/2bf5d756ce45bf0e.
Report an issue: GitHub.