{"record":{"id":"2c38dc0b8d41d923","repo":"risingwavelabs/risingwave","slug":"iceberg-pk-index-writer-expected-barrier-on-rem","errorCode":null,"errorMessage":"iceberg pk-index writer {} expected barrier on remap input, got {msg:?}","messagePattern":"iceberg pk-index writer (.+?) expected barrier on remap input, got (.+?)","errorType":"exception","errorClass":"StreamExecutorError","httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/iceberg_with_pk_index/writer.rs","lineNumber":457,"sourceCode":"    ) {\n        #[for_await]\n        for msg in input {\n            match msg? {\n                Message::Chunk(chunk) =>\n                {\n                    #[for_await]\n                    for chunk in self.process_chunk(chunk) {\n                        yield Message::Chunk(chunk?.into());\n                    }\n                }\n                Message::Watermark(watermark) => {\n                    yield Message::Watermark(watermark);\n                }\n                Message::Barrier(barrier) => {\n                    let msg = next_msg(resolver_input).await?;\n                    let remap_barrier = match msg {\n                        Message::Barrier(b) => b,\n                        _ => bail!(\n                            \"iceberg pk-index writer {} expected barrier on remap input, got {msg:?}\",\n                            self.sink_id\n                        ),\n                    };\n                    self.validate_aligned_barriers(&barrier, &remap_barrier)?;\n                    let begin = self.compaction_begin(&barrier)?;\n                    let epoch = barrier.epoch;\n                    #[for_await]\n                    for msg in self.checkpoint_barrier(barrier) {\n                        yield msg?;\n                    }\n\n                    if let Some(task_id) = begin {\n                        self.mode = WriterInputMode::ResolvingRight {\n                            task_id,\n                            begin_epoch: epoch,\n                        };\n                        return Ok(());","sourceCodeStart":439,"sourceCodeEnd":475,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/iceberg_with_pk_index/writer.rs#L439-L475","documentation":"The Iceberg pk-index sink writer executor keeps two aligned inputs: the main data input and the remap (resolver) input. On every barrier from the main input it expects the remap input to also deliver a barrier so epochs stay aligned; if the remap input delivers a chunk or watermark instead of a barrier at that point, the executor bails with this internal invariant violation.","triggerScenarios":"In execute_normal, a Message::Barrier arrives on the main input, the executor pulls the next message from resolver_input, and that message is a Chunk or Watermark rather than a Barrier (the two upstreams are misaligned).","commonSituations":"A misconfigured or buggy upstream compaction-remap operator that buffers chunks instead of flushing before barriers; a rewritten/patched fragment graph where the remap input is not barrier-aligned with the data input; backpressure causing the remap side to lag behind the data side across epoch boundaries.","solutions":["Verify the stream fragment graph for this iceberg sink (sink_id in the message) still wires the remap/resolver input from the expected compaction remap operator that emits aligned barriers.","Check the upstream operator feeding the remap input for changes that emit chunks/watermarks after the main input's barrier without a preceding barrier.","Reproduce with actor tracing / enable verbose executor logs to see which upstream actor produced the out-of-order message, then fix or roll back the recent change to that operator.","If caused by a mixed-version cluster after upgrade, ensure all nodes run the same RisingWave version so barrier protocol assumptions match."],"exampleFix":"// upstream remap operator must flush buffered chunks before forwarding a barrier\n// before\nbuffer.push(msg); // chunk held across barrier, remap side emits chunk when writer expects barrier\n// after\nif is_barrier(msg) { flush(&mut buffer, out); }\nout.send(msg);","handlingStrategy":"validation","validationCode":"// Before wiring an iceberg sink with pk-index compaction, verify both inputs are barrier-aligned\nfn inputs_barrier_aligned(main_plan: &FragmentPlan, remap_plan: &FragmentPlan) -> bool {\n    main_plan.emits_barrier_every_epoch() && remap_plan.emits_barrier_every_epoch()\n        && main_plan.epoch_alignment == remap_plan.epoch_alignment\n}\nassert!(inputs_barrier_aligned(&main, &remap), \"remap input must be barrier-aligned with data input\");","typeGuard":"fn as_barrier(msg: &Message) -> Option<&Barrier> {\n    match msg { Message::Barrier(b) => Some(b), _ => None }\n}","tryCatchPattern":null,"preventionTips":["Never manually edit the fragment graph feeding an iceberg pk-index sink; use supported DDL only.","Keep all cluster nodes on the same RisingWave version to avoid barrier-protocol skew.","When adding executors upstream, preserve the barrier-first, aligned-barrier contract of the sink inputs."],"tags":["streaming","barrier-alignment","iceberg-sink","internal-invariant"],"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"}