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

  1. Check barrier manager connectivity from the compute node; recover the MV/job.
  2. Inspect meta logs for cancellation or rescheduling of the source actor; retry creation.
  3. 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

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


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)