risingwavelabs/risingwave · error · StreamExecutorError

failed to receive the first barrier, actor_id

Error message

failed to receive the first barrier, actor_id: {:?}, source_id: {:?}

What it means

The batch POSIX FS list source executor fails when the barrier channel from the local barrier manager closes before delivering the first barrier. Sources must receive an initial barrier to establish the epoch before producing data; if the receiver returns None the source cannot start and aborts with this error. It embeds actor_id and source_id for diagnosis.

Solutions

  1. Check compute-node and meta-node connectivity; ensure the barrier manager is running and the actor was not cancelled at startup.
  2. Inspect logs for actor cancellation/failover around the same time; restart the streaming job or resume the MV if it was mid-migration.
  3. If seen in tests, send an initial Barrier into the barrier_receiver before driving the executor.

Example fix

// test setup: before running the source executor
// before
let (tx, rx) = unbounded_channel(); // never sends
// after
let (tx, rx) = unbounded_channel();
tx.send(Barrier::new_test_barrier(1)).unwrap();
Defensive patterns

Strategy: retry

Validate before calling

// check cluster health before assuming source bug
// SELECT * FROM rw_materialized_views WHERE mv_id = <id>; -- state should be CREATED/RUNNING
// verify compute nodes: SELECT * FROM rw_workers WHERE worker_type = 'COMPUTE_NODE';

Try / catch

// on failure, rely on RW recovery; if manually driving:
match barrier_rx.recv().await {
    Some(b) => proceed(b),
    None => Err(recoverable("source never received first barrier")),
}

Prevention

When it happens

Trigger: into_stream() calls barrier_receiver.recv() and the channel is closed/has no senders (barrier manager dropped the sender or the actor is being stopped) so recv() returns None instead of a Barrier.

Common situations: Fragment/actor being cancelled immediately after creation, compute node disconnected from meta during startup, source actor rescheduled during failover, or tests where no barrier was ever sent into the receiver.

Related errors


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

Appendix: source

Thrown at src/stream/src/executor/source/batch_source/batch_posix_fs_list.rs:214

                    } else {
                        files.push((relative_path_str, metadata.len()));
                    }
                }
            }

            Ok(())
        })
    }

    #[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
                )
            })?;

        // 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::BatchPosixFs(batch_posix_fs_properties) = config else {
            unreachable!("BatchPosixFsListExecutor must be used with BatchPosixFs connector")
        };

        yield Message::Barrier(first_barrier);
        let barrier_stream = barrier_to_message_stream(barrier_receiver).boxed();

View on GitHub (pinned to 6469eb736d)