{"record":{"id":"78e004e888269f17","repo":"risingwavelabs/risingwave","slug":"iceberg-pk-index-writer-received-watermark-from","errorCode":null,"errorMessage":"iceberg pk-index writer {} received watermark from replacement input before its initial barrier for task {}","messagePattern":"iceberg pk-index writer (.+?) received watermark 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":537,"sourceCode":"        );\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                }\n            }\n        }\n\n        bail!(\n            \"iceberg pk-index writer {} replacement input closed before its initial barrier for task {}\",","sourceCodeStart":519,"sourceCodeEnd":555,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/iceberg_with_pk_index/writer.rs#L519-L555","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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."],"exampleFix":"// hold watermarks until the replacement input's initial barrier has passed\nmatch msg {\n    Message::Watermark(w) if !initial_barrier_seen => { held.push(w); continue; }\n    other => forward(other),\n}","handlingStrategy":"validation","validationCode":"// Reject replacement-input branches that can forward watermarks before their initial barrier\nfn replacement_branch_emits_early_watermarks(plan: &FragmentPlan) -> bool {\n    plan.head_operators().iter().any(|op| op.emits_watermarks())\n}\nassert!(!replacement_branch_emits_early_watermarks(&plan), \"replacement input must not emit watermarks before its initial barrier\");","typeGuard":"fn is_early_watermark(msg: &Message, initial_barrier_seen: bool) -> bool {\n    matches!(msg, Message::Watermark(_)) && !initial_barrier_seen\n}","tryCatchPattern":null,"preventionTips":["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."],"tags":["streaming","watermark","barrier-alignment","iceberg-sink"],"backgroundTag":"invalid-state-transition","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"}