{"record":{"id":"785699cf71ac0d25","repo":"risingwavelabs/risingwave","slug":"barrier-reader-closed-unexpectedly","errorCode":null,"errorMessage":"barrier reader closed unexpectedly","messagePattern":"barrier reader closed unexpectedly","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/source/mod.rs","lineNumber":72,"sourceCode":"mod source_backfill_state_table;\npub(crate) use source_backfill_state_table::BackfillStateTableHandler;\n\npub mod state_table_handler;\nuse futures_async_stream::try_stream;\nuse risingwave_common::util::retry::exponential_backoff;\nuse tokio::sync::mpsc::UnboundedReceiver;\nuse tokio_retry::strategy::jitter;\n\nuse crate::executor::error::StreamExecutorError;\nuse crate::executor::{Barrier, Message};\n\n/// Receive barriers from barrier manager with the channel, error on channel close.\n#[try_stream(ok = Message, error = StreamExecutorError)]\npub async fn barrier_to_message_stream(mut rx: UnboundedReceiver<Barrier>) {\n    while let Some(barrier) = rx.recv().instrument_await(\"receive_barrier\").await {\n        yield Message::Barrier(barrier);\n    }\n    bail!(\"barrier reader closed unexpectedly\");\n}\n\npub fn get_split_offset_mapping_from_chunk(\n    chunk: &StreamChunk,\n    split_idx: usize,\n    offset_idx: usize,\n) -> Option<HashMap<SplitId, String>> {\n    let mut split_offset_mapping = HashMap::new();\n    // All rows (including those visible or invisible) will be used to update the source offset.\n    for i in 0..chunk.capacity() {\n        let (_, row, _) = chunk.row_at(i);\n        let split_id = row.datum_at(split_idx).unwrap().into_utf8().into();\n        let offset = row.datum_at(offset_idx).unwrap().into_utf8();\n        split_offset_mapping.insert(split_id, offset.to_owned());\n    }\n    Some(split_offset_mapping)\n}\n","sourceCodeStart":54,"sourceCodeEnd":90,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/source/mod.rs#L54-L90","documentation":"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.","triggerScenarios":"All senders to the barrier UnboundedReceiver are dropped while the source executor is running, so the while-let loop exits and bail! is reached.","commonSituations":"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.","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."],"exampleFix":"// before (test)\nlet (tx, rx) = unbounded_channel();\ndrop(tx);\n// after\nlet (tx, rx) = unbounded_channel();\n// keep tx alive, send barriers periodically\nstd::mem::forget(tx); // or hold it in scope","handlingStrategy":"try-catch","validationCode":"// hold a clone/strong reference to the barrier sender while the stream runs\nlet tx_guard = barrier_tx.clone();","typeGuard":null,"tryCatchPattern":"match barrier_to_message_stream(rx).next().await {\n    Err(e) if e.to_string().contains(\"barrier reader closed\") => {\n        // barrier manager gone: trigger recovery/reconnect\n        handle_barrier_loss();\n    }\n    other => propagate(other),\n}","preventionTips":["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"],"tags":["streaming","barrier","channel-closed"],"backgroundTag":"channel-closed-unexpectedly","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"}