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
- Ensure the upstream of the replacement input emits its initial barrier before any data chunks (barrier-first guarantee).
- Diff recent changes to the fragment graph or executors feeding this sink's replacement input and revert the offending change.
- Recreate the sink so the replacement input starts from a clean, barrier-aligned state.
- 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
- Any executor placed at the head of the replacement input must forward its initial barrier before emitting chunks.
- Test custom executors with a barrier-ordering unit test before deploying into an iceberg compaction pipeline.
- Recreate the sink after fragment-graph changes so the replacement input restarts in a clean, aligned state.
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
- iceberg pk-index writer
- iceberg pk-index writer
- iceberg pk-index writer
- iceberg pk-index writer
- build parquet stream for
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)