risingwavelabs/risingwave · error · StreamExecutorError

iceberg pk-index writer

Error message

iceberg pk-index writer {} received watermark on remap input while resolving task {}

What it means

While the writer is in ResolvingRight mode (waiting for the Phase::End compaction barrier on the remap input), the protocol only permits chunks and the final barrier. A watermark arriving on the remap input during this phase breaks the compaction protocol, so the executor bails with this error naming the compaction task.

Solutions

  1. Inspect the fragment graph for this sink's remap input and remove or block the operator that forwards watermarks into it.
  2. If a custom/modified upstream must emit watermarks, drop them before the remap input instead of forwarding.
  3. Check whether the error appeared after an upgrade or plan change and roll back to the last known-good graph definition.
  4. Report the sink_id and task_id to maintainers with the actor/fragment plan if the default pipeline triggers it — it indicates a planner bug.

Example fix

// intercept upstream messages before the remap input
match msg {
    Message::Watermark(_) => continue, // drop
    other => forward(other),
}
Defensive patterns

Strategy: validation

Validate before calling

// Verify the remap input branch contains no watermark-emitting operators before creating the sink
fn remap_branch_emits_watermarks(plan: &FragmentPlan) -> bool {
    plan.operators().iter().any(|op| op.emits_watermarks())
}
assert!(!remap_branch_emits_watermarks(&remap_plan), "remap input must not forward watermarks during compaction");

Type guard

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

Prevention

When it happens

Trigger: In execute_resolving_right, the resolver_input stream yields Message::Watermark while a compaction task (task_id) is being resolved; watermarks are simply not part of the remap input's message contract.

Common situations: A watermarked upstream (e.g., a source or join with watermark emission) is incorrectly wired into the compaction remap input; a plan rewrite or new executor inserted into the remap branch forwards watermarks downstream; a version change adds watermark support to an operator in the remap path.

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/14110b58d0e3a3bd. Report an issue: GitHub.

Appendix: source

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

        }

        *completed = true;
    }

    #[try_stream(ok = Message, error = StreamExecutorError)]
    async fn execute_resolving_right<'a>(
        &'a mut self,
        resolver_input: &'a mut BoxedMessageStream,
        task_id: IcebergCompactionTaskId,
        begin_epoch: EpochPair,
    ) {
        #[for_await]
        for msg in resolver_input {
            match msg? {
                Message::Chunk(chunk) => {
                    self.apply_resolver_chunk(chunk).await?;
                }
                Message::Watermark(_) => bail!(
                    "iceberg pk-index writer {} received watermark on remap input while resolving task {}",
                    self.sink_id,
                    task_id
                ),
                Message::Barrier(barrier) => {
                    self.validate_compaction_barrier(
                        &barrier,
                        task_id,
                        Phase::End,
                        begin_epoch.curr,
                    )?;
                    self.mode = WriterInputMode::AligningReplacementInput { task_id, barrier };
                    return Ok(());
                }
            }
        }

        bail!(

View on GitHub (pinned to 6469eb736d)