risingwavelabs/risingwave · error

barrier reader closed unexpectedly

Error message

barrier reader closed unexpectedly

What it means

barrier_to_message_stream forwards barriers from the barrier manager channel into the executor's message stream. If the channel closes (recv() returns None) the stream deliberately bails with this error, because an executor must never run without barrier flow — a closed channel means the barrier manager is gone.

Solutions

  1. Investigate why the barrier manager stopped sending barriers (node crash, drain, meta disconnect).
  2. Rely on RisingWave recovery: the actor will fail and be recovered; check cluster health if this recurs.
  3. In tests, keep the barrier sender alive (hold `tx`) for the lifetime of the executor stream.

Example fix

// before (test)
let (tx, rx) = unbounded_channel();
drop(tx);
// after
let (tx, rx) = unbounded_channel();
// keep tx alive, send barriers periodically
std::mem::forget(tx); // or hold it in scope
Defensive patterns

Strategy: try-catch

Validate before calling

// hold a clone/strong reference to the barrier sender while the stream runs
let tx_guard = barrier_tx.clone();

Try / catch

match barrier_to_message_stream(rx).next().await {
    Err(e) if e.to_string().contains("barrier reader closed") => {
        // barrier manager gone: trigger recovery/reconnect
        handle_barrier_loss();
    }
    other => propagate(other),
}

Prevention

When it happens

Trigger: All senders to the barrier UnboundedReceiver are dropped while the source executor is running, so the while-let loop exits and bail! is reached.

Common situations: Local barrier manager shutdown, compute node being drained or crashed, failover where the barrier channel is torn down, or test harness dropping the barrier sender.

Related errors


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

Appendix: source

Thrown at src/stream/src/executor/source/mod.rs:72

mod source_backfill_state_table;
pub(crate) use source_backfill_state_table::BackfillStateTableHandler;

pub mod state_table_handler;
use futures_async_stream::try_stream;
use risingwave_common::util::retry::exponential_backoff;
use tokio::sync::mpsc::UnboundedReceiver;
use tokio_retry::strategy::jitter;

use crate::executor::error::StreamExecutorError;
use crate::executor::{Barrier, Message};

/// Receive barriers from barrier manager with the channel, error on channel close.
#[try_stream(ok = Message, error = StreamExecutorError)]
pub async fn barrier_to_message_stream(mut rx: UnboundedReceiver<Barrier>) {
    while let Some(barrier) = rx.recv().instrument_await("receive_barrier").await {
        yield Message::Barrier(barrier);
    }
    bail!("barrier reader closed unexpectedly");
}

pub fn get_split_offset_mapping_from_chunk(
    chunk: &StreamChunk,
    split_idx: usize,
    offset_idx: usize,
) -> Option<HashMap<SplitId, String>> {
    let mut split_offset_mapping = HashMap::new();
    // All rows (including those visible or invisible) will be used to update the source offset.
    for i in 0..chunk.capacity() {
        let (_, row, _) = chunk.row_at(i);
        let split_id = row.datum_at(split_idx).unwrap().into_utf8().into();
        let offset = row.datum_at(offset_idx).unwrap().into_utf8();
        split_offset_mapping.insert(split_id, offset.to_owned());
    }
    Some(split_offset_mapping)
}

View on GitHub (pinned to 6469eb736d)