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
- Inspect the fragment graph for this sink's remap input and remove or block the operator that forwards watermarks into it.
- If a custom/modified upstream must emit watermarks, drop them before the remap input instead of forwarding.
- Check whether the error appeared after an upgrade or plan change and roll back to the last known-good graph definition.
- 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
- Do not place watermark-generating operators (watermark sources, windowed joins) in the remap/resolver branch of an iceberg sink.
- Re-check the plan after upgrades; new versions may change watermark forwarding behavior.
- Recreate the sink rather than mutating its fragment graph while compaction is active.
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
- iceberg pk-index writer
- iceberg pk-index writer
- iceberg pk-index writer
- below watermark check condition eval must return bool array
- build parquet stream for
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)