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
- Investigate why the barrier manager stopped sending barriers (node crash, drain, meta disconnect).
- Rely on RisingWave recovery: the actor will fail and be recovered; check cluster health if this recurs.
- 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
- Never drop the barrier sender while source executors are running
- Ensure graceful shutdown drains executors before tearing down the barrier manager
- In tests, keep tx alive for the whole executor lifetime
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
- end of barrier receiver
- unable to send barrier
- actor exited unexpectedly
- cannot get truncated epoch
- current epoch has exceeded the epoch of the stream that has…
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)