risingwavelabs/risingwave · error

left barrier received while right stream end

Error message

left barrier received while right stream end

What it means

Mirror of the left-end case: during barrier alignment the right input stream ended, and then a barrier arrived from the left input. The Lookup executor treats a barrier after the paired side has terminated as a violation of the termination protocol and fails the actor.

Solutions

  1. Inspect the right upstream's logs to determine why it terminated while the left kept sending barriers.
  2. Re-create the materialized view / sink to get a consistent fresh plan if version skew is suspected.
  3. Align all node versions across the cluster before restarting the stream graph.
  4. If persistent, report with the fragment graph; this indicates an internal scheduling or protocol bug.
Defensive patterns

Strategy: try-catch

Try / catch

// Treat as a job-level failure and restart the stream graph:
if let Err(e) = actor.run().await {
    if e.to_string().contains("left barrier received while right stream end") {
        log::error!("upstream lifecycle mismatch: {e}");
    }
}

Prevention

When it happens

Trigger: Raised inside align_barrier's passthrough loop: after observing a None from the right stream, the loop drains remaining left messages and encounters Message::Barrier from the left input.

Common situations: Asymmetric upstream termination from an actor failure or migration; mixed-version cluster where one side's protocol differs; malformed upstream graph in dev/test harnesses.

Understand the failure class

Background: "Invalid state transition" errors: "status must be X, actually Y", "already rejected/charging/uninstalled", "cannot ... while running" — what they mean when a library rejects your call — this error's family across 31 libraries.

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/78f744d5923562ac. Report an issue: GitHub.

Appendix: source

Thrown at src/stream/src/executor/lookup/sides.rs:158

                    while let Some(msg) = right.next().await {
                        match msg? {
                            w @ Message::Watermark(_) => yield Either::Left(w),
                            c @ Message::Chunk(_) => yield Either::Left(c),
                            Message::Barrier(_) => {
                                bail!("right barrier received while left stream end");
                            }
                        }
                    }
                    break 'outer;
                }
                future::Either::Right((None, _)) => {
                    // right stream end, passthrough left chunks
                    while let Some(msg) = left.next().await {
                        match msg? {
                            w @ Message::Watermark(_) => yield Either::Right(w),
                            c @ Message::Chunk(_) => yield Either::Right(c),
                            Message::Barrier(_) => {
                                bail!("left barrier received while right stream end");
                            }
                        }
                    }
                    break 'outer;
                }
                future::Either::Left((Some(msg), _)) => match msg? {
                    w @ Message::Watermark(_) => yield Either::Left(w),
                    c @ Message::Chunk(_) => yield Either::Left(c),
                    Message::Barrier(b) => {
                        yield Either::Left(Message::Barrier(b.clone()));
                        break 'inner (SideStatus::LeftBarrier, b);
                    }
                },
                future::Either::Right((Some(msg), _)) => match msg? {
                    w @ Message::Watermark(_) => yield Either::Right(w),
                    c @ Message::Chunk(_) => yield Either::Right(c),
                    Message::Barrier(b) => {
                        yield Either::Right(Message::Barrier(b.clone()));

View on GitHub (pinned to 6469eb736d)