{"record":{"id":"ebfe251afeab4c64","repo":"risingwavelabs/risingwave","slug":"iceberg-pk-index-writer-received-a-chunk-from-r","errorCode":null,"errorMessage":"iceberg pk-index writer {} received a chunk from replacement input before its initial barrier for task {}","messagePattern":"iceberg pk-index writer (.+?) received a chunk from replacement input before its initial barrier for task (.+?)","errorType":"exception","errorClass":"StreamExecutorError","httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/iceberg_with_pk_index/writer.rs","lineNumber":532,"sourceCode":"\n        bail!(\n            \"iceberg pk-index writer {} resolver input closed before switch-to-input for task {}\",\n            self.sink_id,\n            task_id\n        );\n    }\n\n    #[try_stream(ok = Message, error = StreamExecutorError)]\n    async fn execute_aligning_replacement_input<'a>(\n        &'a mut self,\n        input: &'a mut BoxedMessageStream,\n        task_id: IcebergCompactionTaskId,\n        expected: Barrier,\n    ) {\n        #[for_await]\n        for msg in input {\n            match msg? {\n                Message::Chunk(_) => bail!(\n                    \"iceberg pk-index writer {} received a chunk from replacement input before its initial barrier for task {}\",\n                    self.sink_id,\n                    task_id\n                ),\n                Message::Watermark(_) => bail!(\n                    \"iceberg pk-index writer {} received watermark from replacement input before its initial barrier for task {}\",\n                    self.sink_id,\n                    task_id\n                ),\n                Message::Barrier(barrier) => {\n                    self.validate_aligned_barriers(&barrier, &expected)?;\n                    #[for_await]\n                    for msg in self.checkpoint_barrier(barrier) {\n                        yield msg?;\n                    }\n                    self.mode = WriterInputMode::Normal;\n                    return Ok(());\n                }","sourceCodeStart":514,"sourceCodeEnd":550,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/iceberg_with_pk_index/writer.rs#L514-L550","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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."],"exampleFix":"// replacement-input executor must not emit data before the first barrier\n// before\nif let Some(chunk) = pending.pop() { out.send(Message::Chunk(chunk)); }\n// after\nif !initial_barrier_sent { out.send(Message::Barrier(first_barrier)); initial_barrier_sent = true; }","handlingStrategy":"validation","validationCode":"// Ensure the replacement-input executor implements the barrier-first guarantee before enabling compaction\nfn replacement_input_barrier_first(op: &dyn Executor) -> bool {\n    op.first_message_kind() == Some(MessageKind::Barrier)\n}\nassert!(replacement_input_barrier_first(&replacement_executor), \"replacement input must emit its initial barrier before any data\");","typeGuard":"fn first_message_is_barrier(msg: &Message) -> bool {\n    matches!(msg, Message::Barrier(_))\n}","tryCatchPattern":null,"preventionTips":["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."],"tags":["streaming","barrier-alignment","iceberg-sink","compaction"],"backgroundTag":"internal-invariant-violation","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}