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

  1. Compare the logged expected_prev with the barrier's epoch.prev to find where the chain diverged.
  2. Restart the fragment from the last consistent checkpoint so epochs realign.
  3. Check meta's compaction scheduler for phase transitions issued on wrong barriers.
  4. Verify no configuration change (checkpoint interval, actor layout) occurred mid-compaction.
  5. 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

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


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/2bf5d756ce45bf0e. Report an issue: GitHub.