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
- Remove or gate the watermark-emitting operator in the replacement input branch of this sink's fragment graph.
- Drop watermarks until after the initial barrier is forwarded if a custom upstream must emit them.
- Check for a recent plan change or upgrade that introduced watermark forwarding on this path and roll back.
- 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
- Keep watermark-emitting operators out of the replacement input branch, or buffer watermarks until after the initial barrier.
- Re-validate the plan after version upgrades that change watermark propagation.
- Report any unmodified-pipeline occurrence to maintainers with sink_id, task_id, and the actor graph — this is a protocol violation, not a user-data issue.
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
- iceberg pk-index writer
- iceberg pk-index writer
- iceberg pk-index writer
- below watermark check condition eval must return bool array
- Expected at most 1 clean_watermark_index per table, got
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)