{"record":{"id":"14110b58d0e3a3bd","repo":"risingwavelabs/risingwave","slug":"iceberg-pk-index-writer-received-watermark-on-r","errorCode":null,"errorMessage":"iceberg pk-index writer {} received watermark on remap input while resolving task {}","messagePattern":"iceberg pk-index writer (.+?) received watermark on remap input while resolving task (.+?)","errorType":"exception","errorClass":"StreamExecutorError","httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/iceberg_with_pk_index/writer.rs","lineNumber":497,"sourceCode":"        }\n\n        *completed = true;\n    }\n\n    #[try_stream(ok = Message, error = StreamExecutorError)]\n    async fn execute_resolving_right<'a>(\n        &'a mut self,\n        resolver_input: &'a mut BoxedMessageStream,\n        task_id: IcebergCompactionTaskId,\n        begin_epoch: EpochPair,\n    ) {\n        #[for_await]\n        for msg in resolver_input {\n            match msg? {\n                Message::Chunk(chunk) => {\n                    self.apply_resolver_chunk(chunk).await?;\n                }\n                Message::Watermark(_) => bail!(\n                    \"iceberg pk-index writer {} received watermark on remap input while resolving task {}\",\n                    self.sink_id,\n                    task_id\n                ),\n                Message::Barrier(barrier) => {\n                    self.validate_compaction_barrier(\n                        &barrier,\n                        task_id,\n                        Phase::End,\n                        begin_epoch.curr,\n                    )?;\n                    self.mode = WriterInputMode::AligningReplacementInput { task_id, barrier };\n                    return Ok(());\n                }\n            }\n        }\n\n        bail!(","sourceCodeStart":479,"sourceCodeEnd":515,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/iceberg_with_pk_index/writer.rs#L479-L515","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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."],"exampleFix":"// intercept upstream messages before the remap input\nmatch msg {\n    Message::Watermark(_) => continue, // drop\n    other => forward(other),\n}","handlingStrategy":"validation","validationCode":"// Verify the remap input branch contains no watermark-emitting operators before creating the sink\nfn remap_branch_emits_watermarks(plan: &FragmentPlan) -> bool {\n    plan.operators().iter().any(|op| op.emits_watermarks())\n}\nassert!(!remap_branch_emits_watermarks(&remap_plan), \"remap input must not forward watermarks during compaction\");","typeGuard":"fn is_watermark(msg: &Message) -> bool {\n    matches!(msg, Message::Watermark(_))\n}","tryCatchPattern":null,"preventionTips":["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."],"tags":["streaming","watermark","iceberg-sink","compaction"],"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"}