risingwavelabs/risingwave · error
failed to receive the first barrier, actor_id
Error message
failed to receive the first barrier, actor_id: {:?}, source_id: {:?} What it means
The main stream source executor must receive the first barrier before producing any data so it can anchor the initial epoch. If the barrier channel closes before the first barrier arrives, execute_inner fails with this error containing actor_id and source_id.
Solutions
- Check meta/compute connectivity and barrier manager health; recover the MV or job.
- Look for actor cancellation/rescheduling logs at the same timestamp; recreate the source if the fragment is stuck.
- In tests, send an initial Barrier into the channel before driving the executor.
Defensive patterns
Strategy: retry
Validate before calling
// preflight: check source fragment running state // SELECT * FROM rw_sources WHERE id = <source_id>; -- expect RUNNING/CREATED
Try / catch
let first = match barrier_rx.recv().await {
Some(b) => b,
None => return Err(recover_source("source executor: no first barrier")),
};
let first_epoch = first.epoch; Prevention
- Keep meta/compute healthy during source startup
- Alert on frequent actor recreation or barrier delivery gaps
- In tests, prime the barrier channel with an initial barrier
When it happens
Trigger: execute_inner() calls barrier_receiver.recv() and gets None because the barrier manager's sender side was dropped or never connected to this actor.
Common situations: Actor cancelled at startup, meta/compute disconnect, failover racing source initialization, or tests where no barrier was ever sent.
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/6b115d5e3f43db6d.
Report an issue: GitHub.
Appendix: source
Thrown at src/stream/src/executor/source/source_executor.rs:623
);
// Mark as reported to prevent any future reports, even if offset changes
*must_report_cdc_offset_once = false;
}
}
/// A source executor with a stream source receives:
/// 1. Barrier messages
/// 2. Data from external source
/// and acts accordingly.
#[try_stream(ok = Message, error = StreamExecutorError)]
async fn execute_inner(mut self) {
let mut barrier_receiver = self.barrier_receiver.take().unwrap();
let first_barrier = barrier_receiver
.recv()
.instrument_await("source_recv_first_barrier")
.await
.ok_or_else(|| {
anyhow!(
"failed to receive the first barrier, actor_id: {:?}, source_id: {:?}",
self.actor_ctx.id,
self.stream_source_core.source_id
)
})?;
let first_epoch = first_barrier.epoch;
// must_report_cdc_offset is true if and only if the source is a CDC source.
// must_wait_cdc_offset_before_report is true if and only if the source is a MySQL or SQL Server CDC source.
let (mut boot_state, mut must_report_cdc_offset_once, must_wait_cdc_offset_before_report) =
if let Some(splits) = first_barrier.initial_split_assignment(self.actor_ctx.id) {
// CDC source must reach this branch.
tracing::debug!(?splits, "boot with splits");
// Skip report for non-CDC.
let must_report_cdc_offset_once = splits.iter().any(|split| split.is_cdc_split());
// Only for MySQL and SQL Server CDC, we need to wait for the offset to be non-empty before reporting.
let must_wait_cdc_offset_before_report = must_report_cdc_offset_once
&& splits.iter().any(|split| {
matches!(split, SplitImpl::MysqlCdc(_) | SplitImpl::SqlServerCdc(_))View on GitHub (pinned to 6469eb736d)