{"record":{"id":"9eb9d283a779c120","repo":"risingwavelabs/risingwave","slug":"failed-to-receive-the-first-barrier-actor-id-9eb9d2","errorCode":null,"errorMessage":"failed to receive the first barrier, actor_id: {:?} with no stream source","messagePattern":"failed to receive the first barrier, actor_id: (.+?) with no stream source","errorType":"exception","errorClass":"StreamExecutorError","httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/source/dummy_source_executor.rs","lineNumber":49,"sourceCode":"impl DummySourceExecutor {\n    pub fn new(actor_ctx: ActorContextRef, barrier_receiver: UnboundedReceiver<Barrier>) -> Self {\n        Self {\n            actor_ctx,\n            barrier_receiver: Some(barrier_receiver),\n        }\n    }\n\n    /// A dummy source executor only receives barrier messages and sends them to\n    /// the downstream executor.\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 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: {:?} with no stream source\",\n                    self.actor_ctx.id\n                )\n            })?;\n        yield Message::Barrier(barrier);\n\n        while let Some(barrier) = barrier_receiver.recv().await {\n            yield Message::Barrier(barrier);\n        }\n    }\n}\n\nimpl Execute for DummySourceExecutor {\n    fn execute(self: Box<Self>) -> BoxedMessageStream {\n        self.execute_inner().boxed()\n    }\n}\n","sourceCodeStart":31,"sourceCodeEnd":67,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/source/dummy_source_executor.rs#L31-L67","documentation":"The dummy source executor (a source with no actual stream source, used e.g. for DML-only or CDC-tail setups) must still receive a first barrier before emitting anything. If the barrier channel closes before the first barrier arrives, execute_inner fails with this error including the actor_id.","triggerScenarios":"execute_inner() calls barrier_receiver.recv(); recv() returns None because all senders (the local barrier manager's) were dropped before any barrier was sent.","commonSituations":"Actor cancelled during startup, meta/compute disconnect during recovery, or unit tests constructing DummySourceExecutor without wiring a barrier sender.","solutions":["Ensure the barrier manager is sending barriers to this actor; check meta/compute connectivity.","Look for concurrent actor cancellation or failover in logs; retry the job/MV recovery.","In tests, send an initial barrier into the receiver before polling the executor stream."],"exampleFix":"// before\nlet mut exec = DummySourceExecutor::new(..., rx, ...);\n// after (test)\nbarrier_tx.send(Barrier::new_test_barrier(1)).unwrap();\nlet mut exec = DummySourceExecutor::new(..., rx, ...);","handlingStrategy":"retry","validationCode":"// ensure barrier manager is up before starting actors\n// check compute node logs for 'barrier manager' readiness","typeGuard":null,"tryCatchPattern":"match barrier_rx.recv().await {\n    Some(b) => Ok(b),\n    None => Err(RecoveryNeeded(\"dummy source: no first barrier\")),\n}","preventionTips":["Keep barrier senders alive for the actor's lifetime","Avoid racing actor cancellation with executor startup","Send a first barrier in test harnesses before polling"],"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"}