risingwavelabs/risingwave · error · StreamExecutorError

failed to receive the first barrier, actor_id

Error message

failed to receive the first barrier, actor_id: {:?} with no stream source

What it means

The dummy source executor (a source with no actual stream source, used e.g. for DML-only or CDC-tail setups) must still receive a first barrier before emitting anything. If the barrier channel closes before the first barrier arrives, execute_inner fails with this error including the actor_id.

Solutions

  1. Ensure the barrier manager is sending barriers to this actor; check meta/compute connectivity.
  2. Look for concurrent actor cancellation or failover in logs; retry the job/MV recovery.
  3. In tests, send an initial barrier into the receiver before polling the executor stream.

Example fix

// before
let mut exec = DummySourceExecutor::new(..., rx, ...);
// after (test)
barrier_tx.send(Barrier::new_test_barrier(1)).unwrap();
let mut exec = DummySourceExecutor::new(..., rx, ...);
Defensive patterns

Strategy: retry

Validate before calling

// ensure barrier manager is up before starting actors
// check compute node logs for 'barrier manager' readiness

Try / catch

match barrier_rx.recv().await {
    Some(b) => Ok(b),
    None => Err(RecoveryNeeded("dummy source: no first barrier")),
}

Prevention

When it happens

Trigger: execute_inner() calls barrier_receiver.recv(); recv() returns None because all senders (the local barrier manager's) were dropped before any barrier was sent.

Common situations: Actor cancelled during startup, meta/compute disconnect during recovery, or unit tests constructing DummySourceExecutor without wiring a barrier sender.

Related errors


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

Appendix: source

Thrown at src/stream/src/executor/source/dummy_source_executor.rs:49

impl DummySourceExecutor {
    pub fn new(actor_ctx: ActorContextRef, barrier_receiver: UnboundedReceiver<Barrier>) -> Self {
        Self {
            actor_ctx,
            barrier_receiver: Some(barrier_receiver),
        }
    }

    /// A dummy source executor only receives barrier messages and sends them to
    /// the downstream executor.
    #[try_stream(ok = Message, error = StreamExecutorError)]
    async fn execute_inner(mut self) {
        let mut barrier_receiver = self.barrier_receiver.take().unwrap();
        let barrier = barrier_receiver
            .recv()
            .instrument_await("source_recv_first_barrier")
            .await
            .ok_or_else(|| {
                anyhow!(
                    "failed to receive the first barrier, actor_id: {:?} with no stream source",
                    self.actor_ctx.id
                )
            })?;
        yield Message::Barrier(barrier);

        while let Some(barrier) = barrier_receiver.recv().await {
            yield Message::Barrier(barrier);
        }
    }
}

impl Execute for DummySourceExecutor {
    fn execute(self: Box<Self>) -> BoxedMessageStream {
        self.execute_inner().boxed()
    }
}

View on GitHub (pinned to 6469eb736d)