{"record":{"id":"45770546f7b54554","repo":"risingwavelabs/risingwave","slug":"failed-to-receive-the-first-barrier-actor-id-457705","errorCode":null,"errorMessage":"failed to receive the first barrier, actor_id: {:?}, source_id: {:?}","messagePattern":"failed to receive the first barrier, actor_id: (.+?), source_id: (.+?)","errorType":"exception","errorClass":"StreamExecutorError","httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/source/batch_source/batch_posix_fs_list.rs","lineNumber":214,"sourceCode":"                    } else {\n                        files.push((relative_path_str, metadata.len()));\n                    }\n                }\n            }\n\n            Ok(())\n        })\n    }\n\n    #[try_stream(ok = Message, error = StreamExecutorError)]\n    async fn into_stream(mut self) {\n        let mut barrier_receiver = self.barrier_receiver.take().unwrap();\n        let first_barrier = barrier_receiver\n            .recv()\n            .instrument_await(\"source_recv_first_barrier\")\n            .await\n            .ok_or_else(|| {\n                anyhow!(\n                    \"failed to receive the first barrier, actor_id: {:?}, source_id: {:?}\",\n                    self.actor_ctx.id,\n                    self.stream_source_core.source_id\n                )\n            })?;\n\n        // Build source description from the builder.\n        let source_desc_builder: SourceDescBuilder =\n            self.stream_source_core.source_desc_builder.take().unwrap();\n\n        let properties = source_desc_builder.with_properties();\n        let config = ConnectorProperties::extract(properties, false)?;\n        let ConnectorProperties::BatchPosixFs(batch_posix_fs_properties) = config else {\n            unreachable!(\"BatchPosixFsListExecutor must be used with BatchPosixFs connector\")\n        };\n\n        yield Message::Barrier(first_barrier);\n        let barrier_stream = barrier_to_message_stream(barrier_receiver).boxed();","sourceCodeStart":196,"sourceCodeEnd":232,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/source/batch_source/batch_posix_fs_list.rs#L196-L232","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","solutions":["Check compute-node and meta-node connectivity; ensure the barrier manager is running and the actor was not cancelled at startup.","Inspect logs for actor cancellation/failover around the same time; restart the streaming job or resume the MV if it was mid-migration.","If seen in tests, send an initial Barrier into the barrier_receiver before driving the executor."],"exampleFix":"// test setup: before running the source executor\n// before\nlet (tx, rx) = unbounded_channel(); // never sends\n// after\nlet (tx, rx) = unbounded_channel();\ntx.send(Barrier::new_test_barrier(1)).unwrap();","handlingStrategy":"retry","validationCode":"// check cluster health before assuming source bug\n// SELECT * FROM rw_materialized_views WHERE mv_id = <id>; -- state should be CREATED/RUNNING\n// verify compute nodes: SELECT * FROM rw_workers WHERE worker_type = 'COMPUTE_NODE';","typeGuard":null,"tryCatchPattern":"// on failure, rely on RW recovery; if manually driving:\nmatch barrier_rx.recv().await {\n    Some(b) => proceed(b),\n    None => Err(recoverable(\"source never received first barrier\")),\n}","preventionTips":["Keep meta and compute nodes connected during source creation","Avoid cancelling fragments immediately after creation","In tests, always send an initial barrier before polling the executor"],"tags":["streaming","barrier","source-executor"],"backgroundTag":"channel-closed-before-first-barrier","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}