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
- Ensure the barrier manager is sending barriers to this actor; check meta/compute connectivity.
- Look for concurrent actor cancellation or failover in logs; retry the job/MV recovery.
- 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
- Keep barrier senders alive for the actor's lifetime
- Avoid racing actor cancellation with executor startup
- Send a first barrier in test harnesses before polling
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
- failed to receive the first barrier, actor_id
- failed to receive the first barrier, actor_id
- actor exited unexpectedly
- barrier reader closed unexpectedly
- 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/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)