risingwavelabs/risingwave · error · StreamExecutorError

iceberg pk-index writer

Error message

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

What it means

Same alignment phase as the chunk variant: while waiting for the replacement input's initial barrier for an Iceberg compaction task, the writer instead received a Watermark. Watermarks are not permitted from the replacement input before its initial barrier, so the executor bails with the task id in the message.

Solutions

  1. Remove or gate the watermark-emitting operator in the replacement input branch of this sink's fragment graph.
  2. Drop watermarks until after the initial barrier is forwarded if a custom upstream must emit them.
  3. Check for a recent plan change or upgrade that introduced watermark forwarding on this path and roll back.
  4. If seen on an unmodified pipeline, gather sink_id/task_id and the actor graph and report it as a planner/executor bug.

Example fix

// hold watermarks until the replacement input's initial barrier has passed
match msg {
    Message::Watermark(w) if !initial_barrier_seen => { held.push(w); continue; }
    other => forward(other),
}
Defensive patterns

Strategy: validation

Validate before calling

// Reject replacement-input branches that can forward watermarks before their initial barrier
fn replacement_branch_emits_early_watermarks(plan: &FragmentPlan) -> bool {
    plan.head_operators().iter().any(|op| op.emits_watermarks())
}
assert!(!replacement_branch_emits_early_watermarks(&plan), "replacement input must not emit watermarks before its initial barrier");

Type guard

fn is_early_watermark(msg: &Message, initial_barrier_seen: bool) -> bool {
    matches!(msg, Message::Watermark(_)) && !initial_barrier_seen
}

Prevention

When it happens

Trigger: In execute_aligning_replacement_input, the first message pulled from the replacement main input is Message::Watermark instead of the expected initial Message::Barrier for task_id.

Common situations: A watermark-generating operator (source with watermark column, time-window join) sits at the head of the replacement input branch and forwards watermarks before barriers; a plan rewrite introduced watermark emission into the replacement path; version skew between nodes with differing watermark-forwarding behavior.

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/78e004e888269f17. Report an issue: GitHub.

Appendix: source

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

        );
    }

    #[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(());
                }
            }
        }

        bail!(
            "iceberg pk-index writer {} replacement input closed before its initial barrier for task {}",

View on GitHub (pinned to 6469eb736d)