risingwavelabs/risingwave · error · StreamExecutorError

iceberg pk-index writer

Error message

iceberg pk-index writer {} received a chunk from replacement input before its initial barrier for task {}

What it means

After a compaction task's end barrier is seen on the remap input, the writer switches data sources and must first see the replacement input's initial barrier before any data. Receiving a Chunk from the replacement input before that initial barrier breaks the alignment protocol for the task, so the executor bails.

Solutions

  1. Ensure the upstream of the replacement input emits its initial barrier before any data chunks (barrier-first guarantee).
  2. Diff recent changes to the fragment graph or executors feeding this sink's replacement input and revert the offending change.
  3. Recreate the sink so the replacement input starts from a clean, barrier-aligned state.
  4. If reproducible on an unmodified pipeline, report with sink_id, task_id, and the actor plan — this indicates a core bug.

Example fix

// replacement-input executor must not emit data before the first barrier
// before
if let Some(chunk) = pending.pop() { out.send(Message::Chunk(chunk)); }
// after
if !initial_barrier_sent { out.send(Message::Barrier(first_barrier)); initial_barrier_sent = true; }
Defensive patterns

Strategy: validation

Validate before calling

// Ensure the replacement-input executor implements the barrier-first guarantee before enabling compaction
fn replacement_input_barrier_first(op: &dyn Executor) -> bool {
    op.first_message_kind() == Some(MessageKind::Barrier)
}
assert!(replacement_input_barrier_first(&replacement_executor), "replacement input must emit its initial barrier before any data");

Type guard

fn first_message_is_barrier(msg: &Message) -> bool {
    matches!(msg, Message::Barrier(_))
}

Prevention

When it happens

Trigger: In execute_aligning_replacement_input, the very first non-error message pulled from the (replacement) main input is Message::Chunk instead of Message::Barrier for the given compaction task_id.

Common situations: The operator upstream of the replacement input started emitting data before its first barrier (e.g., a freshly created/rewritten fragment that streams chunks immediately); a buggy or patched executor in the replacement path; snapshot/backfill logic injecting chunks ahead of barrier alignment.

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


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

Appendix: source

Thrown at src/stream/src/executor/iceberg_with_pk_index/writer.rs:532

        bail!(
            "iceberg pk-index writer {} resolver input closed before switch-to-input for task {}",
            self.sink_id,
            task_id
        );
    }

    #[try_stream(ok = Message, error = StreamExecutorError)]
    async fn execute_aligning_replacement_input<'a>(
        &'a mut self,
        input: &'a mut BoxedMessageStream,
        task_id: IcebergCompactionTaskId,
        expected: Barrier,
    ) {
        #[for_await]
        for msg in input {
            match msg? {
                Message::Chunk(_) => bail!(
                    "iceberg pk-index writer {} received a chunk from replacement input before its initial barrier for task {}",
                    self.sink_id,
                    task_id
                ),
                Message::Watermark(_) => bail!(
                    "iceberg pk-index writer {} received watermark from replacement input before its initial barrier for task {}",
                    self.sink_id,
                    task_id
                ),
                Message::Barrier(barrier) => {
                    self.validate_aligned_barriers(&barrier, &expected)?;
                    #[for_await]
                    for msg in self.checkpoint_barrier(barrier) {
                        yield msg?;
                    }
                    self.mode = WriterInputMode::Normal;
                    return Ok(());
                }

View on GitHub (pinned to 6469eb736d)