risingwavelabs/risingwave · error

end of barrier receiver

Error message

end of barrier receiver

What it means

The unbounded barrier receiver returned None, meaning the barrier channel sender was dropped. The barrier stream is the executor's heartbeat and control plane, so a closed channel means the actor's barrier input is gone and the executor cannot continue.

Solutions

  1. Look at meta/actor logs preceding the error to find why the barrier sender was dropped (actor failure, cancel, shutdown).
  2. Ensure the barrier channel is kept alive for the executor's full lifetime after rebuild on reschedule.
  3. If it happens without any shutdown/failure, capture a trace and file an issue.
Defensive patterns

Strategy: try-catch

Try / catch

let barrier = receive_next_barrier(&mut barrier_rx).await.inspect_err(|_| warn!("barrier channel closed; actor cannot continue"))?;

Prevention

When it happens

Trigger: All `Barrier` senders dropped while `receive_next_barrier` polls — typically the actor handle/multi-channel honoree (actor context) is dropped during shutdown, failure, or cancellation.

Common situations: Cluster shutdown or actor failure mid-stream; meta node terminating the actor; bugs in barrier channel lifecycle during actor reschedule.

Related errors


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

Appendix: source

Thrown at src/stream/src/executor/backfill/snapshot_backfill/utils.rs:26

//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

use anyhow::anyhow;
use tokio::sync::mpsc::UnboundedReceiver;

use crate::executor::{Barrier, StreamExecutorResult};

pub(super) async fn receive_next_barrier(
    barrier_rx: &mut UnboundedReceiver<Barrier>,
) -> StreamExecutorResult<Barrier> {
    Ok(barrier_rx
        .recv()
        .await
        .ok_or_else(|| anyhow!("end of barrier receiver"))?)
}

View on GitHub (pinned to 6469eb736d)