risingwavelabs/risingwave · error · StreamExecutorError
iceberg pk-index writer
Error message
iceberg pk-index writer {} expected barrier on remap input, got {msg:?} What it means
The Iceberg pk-index sink writer executor keeps two aligned inputs: the main data input and the remap (resolver) input. On every barrier from the main input it expects the remap input to also deliver a barrier so epochs stay aligned; if the remap input delivers a chunk or watermark instead of a barrier at that point, the executor bails with this internal invariant violation.
Solutions
- Verify the stream fragment graph for this iceberg sink (sink_id in the message) still wires the remap/resolver input from the expected compaction remap operator that emits aligned barriers.
- Check the upstream operator feeding the remap input for changes that emit chunks/watermarks after the main input's barrier without a preceding barrier.
- Reproduce with actor tracing / enable verbose executor logs to see which upstream actor produced the out-of-order message, then fix or roll back the recent change to that operator.
- If caused by a mixed-version cluster after upgrade, ensure all nodes run the same RisingWave version so barrier protocol assumptions match.
Example fix
// upstream remap operator must flush buffered chunks before forwarding a barrier
// before
buffer.push(msg); // chunk held across barrier, remap side emits chunk when writer expects barrier
// after
if is_barrier(msg) { flush(&mut buffer, out); }
out.send(msg); Defensive patterns
Strategy: validation
Validate before calling
// Before wiring an iceberg sink with pk-index compaction, verify both inputs are barrier-aligned
fn inputs_barrier_aligned(main_plan: &FragmentPlan, remap_plan: &FragmentPlan) -> bool {
main_plan.emits_barrier_every_epoch() && remap_plan.emits_barrier_every_epoch()
&& main_plan.epoch_alignment == remap_plan.epoch_alignment
}
assert!(inputs_barrier_aligned(&main, &remap), "remap input must be barrier-aligned with data input"); Type guard
fn as_barrier(msg: &Message) -> Option<&Barrier> {
match msg { Message::Barrier(b) => Some(b), _ => None }
} Prevention
- Never manually edit the fragment graph feeding an iceberg pk-index sink; use supported DDL only.
- Keep all cluster nodes on the same RisingWave version to avoid barrier-protocol skew.
- When adding executors upstream, preserve the barrier-first, aligned-barrier contract of the sink inputs.
When it happens
Trigger: In execute_normal, a Message::Barrier arrives on the main input, the executor pulls the next message from resolver_input, and that message is a Chunk or Watermark rather than a Barrier (the two upstreams are misaligned).
Common situations: A misconfigured or buggy upstream compaction-remap operator that buffers chunks instead of flushing before barriers; a rewritten/patched fragment graph where the remap input is not barrier-aligned with the data input; backpressure causing the remap side to lag behind the data side across epoch boundaries.
Understand the failure class
Background: "This is a bug, please report it": internal invariant violations, unreachable panics, and SNH errors explained — this error's family across 47 libraries.
Related errors
- iceberg pk-index writer
- iceberg pk-index writer
- iceberg pk-index writer
- iceberg pk-index writer
- iceberg pk-index writer
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/2c38dc0b8d41d923.
Report an issue: GitHub.
Appendix: source
Thrown at src/stream/src/executor/iceberg_with_pk_index/writer.rs:457
) {
#[for_await]
for msg in input {
match msg? {
Message::Chunk(chunk) =>
{
#[for_await]
for chunk in self.process_chunk(chunk) {
yield Message::Chunk(chunk?.into());
}
}
Message::Watermark(watermark) => {
yield Message::Watermark(watermark);
}
Message::Barrier(barrier) => {
let msg = next_msg(resolver_input).await?;
let remap_barrier = match msg {
Message::Barrier(b) => b,
_ => bail!(
"iceberg pk-index writer {} expected barrier on remap input, got {msg:?}",
self.sink_id
),
};
self.validate_aligned_barriers(&barrier, &remap_barrier)?;
let begin = self.compaction_begin(&barrier)?;
let epoch = barrier.epoch;
#[for_await]
for msg in self.checkpoint_barrier(barrier) {
yield msg?;
}
if let Some(task_id) = begin {
self.mode = WriterInputMode::ResolvingRight {
task_id,
begin_epoch: epoch,
};
return Ok(());View on GitHub (pinned to 6469eb736d)