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 Iceberg list executor must receive the first barrier to learn the initial epoch before listing table snapshots. If the barrier channel is closed with no message, into_stream returns this error with actor_id and source_id, aborting the source.
Solutions
- Check barrier manager connectivity from the compute node; recover the MV/job.
- Inspect meta logs for cancellation or rescheduling of the source actor; retry creation.
- In tests, send an initial Barrier into the channel before polling.
Defensive patterns
Strategy: retry
Validate before calling
// preflight: ensure compute node has barrier manager connection before creating Iceberg sources
Try / catch
match barrier_rx.recv().await {
Some(b) => Ok(b.epoch),
None => Err(recover_with_backoff("iceberg list: no first barrier")),
} Prevention
- Keep the cluster stable during CREATE SOURCE for Iceberg tables
- Monitor barrier manager health and actor recovery events
- In tests, send an initial barrier before polling the executor
When it happens
Trigger: into_stream() awaits barrier_receiver.recv() and gets None (channel closed / no senders) instead of the initial Barrier.
Common situations: Iceberg source actor cancelled at startup, compute node lost barrier connection to meta, or recovery/failover racing executor initialization.
Related errors
- iceberg pk-index writer
- actor exited unexpectedly
- barrier reader closed unexpectedly
- compaction resolver PK column missing column_desc
- compaction resolver sink
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/f086ffd805206a1a.
Report an issue: GitHub.
Appendix: source
Thrown at src/stream/src/executor/source/iceberg_list_executor.rs:99
stream_source_core,
downstream_columns,
metrics,
barrier_receiver: Some(barrier_receiver),
system_params,
rate_limit_rps,
streaming_config,
}
}
#[try_stream(ok = Message, error = StreamExecutorError)]
async fn into_stream(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;
// Build source description from the builder.
let source_desc_builder: SourceDescBuilder =
self.stream_source_core.source_desc_builder.take().unwrap();
let properties = source_desc_builder.with_properties();
let config = ConnectorProperties::extract(properties, false)?;
let ConnectorProperties::Iceberg(iceberg_properties) = config else {
unreachable!()
};
let scan_projection =View on GitHub (pinned to 6469eb736d)