{"record":{"id":"6b115d5e3f43db6d","repo":"risingwavelabs/risingwave","slug":"failed-to-receive-the-first-barrier-actor-id-6b115d","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":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/source/source_executor.rs","lineNumber":623,"sourceCode":"            );\n            // Mark as reported to prevent any future reports, even if offset changes\n            *must_report_cdc_offset_once = false;\n        }\n    }\n\n    /// A source executor with a stream source receives:\n    /// 1. Barrier messages\n    /// 2. Data from external source\n    /// and acts accordingly.\n    #[try_stream(ok = Message, error = StreamExecutorError)]\n    async fn execute_inner(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        let first_epoch = first_barrier.epoch;\n        // must_report_cdc_offset is true if and only if the source is a CDC source.\n        // must_wait_cdc_offset_before_report is true if and only if the source is a MySQL or SQL Server CDC source.\n        let (mut boot_state, mut must_report_cdc_offset_once, must_wait_cdc_offset_before_report) =\n            if let Some(splits) = first_barrier.initial_split_assignment(self.actor_ctx.id) {\n                // CDC source must reach this branch.\n                tracing::debug!(?splits, \"boot with splits\");\n                // Skip report for non-CDC.\n                let must_report_cdc_offset_once = splits.iter().any(|split| split.is_cdc_split());\n                // Only for MySQL and SQL Server CDC, we need to wait for the offset to be non-empty before reporting.\n                let must_wait_cdc_offset_before_report = must_report_cdc_offset_once\n                    && splits.iter().any(|split| {\n                        matches!(split, SplitImpl::MysqlCdc(_) | SplitImpl::SqlServerCdc(_))","sourceCodeStart":605,"sourceCodeEnd":641,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/source/source_executor.rs#L605-L641","documentation":"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.","triggerScenarios":"execute_inner() calls barrier_receiver.recv() and gets None because the barrier manager's sender side was dropped or never connected to this actor.","commonSituations":"Actor cancelled at startup, meta/compute disconnect, failover racing source initialization, or tests where no barrier was ever sent.","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."],"exampleFix":null,"handlingStrategy":"retry","validationCode":"// preflight: check source fragment running state\n// SELECT * FROM rw_sources WHERE id = <source_id>; -- expect RUNNING/CREATED","typeGuard":null,"tryCatchPattern":"let first = match barrier_rx.recv().await {\n    Some(b) => b,\n    None => return Err(recover_source(\"source executor: no first barrier\")),\n};\nlet first_epoch = first.epoch;","preventionTips":["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"],"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"}