{"record":{"id":"f086ffd805206a1a","repo":"risingwavelabs/risingwave","slug":"failed-to-receive-the-first-barrier-actor-id-f086ff","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/iceberg_list_executor.rs","lineNumber":99,"sourceCode":"            stream_source_core,\n            downstream_columns,\n            metrics,\n            barrier_receiver: Some(barrier_receiver),\n            system_params,\n            rate_limit_rps,\n            streaming_config,\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        let first_epoch = first_barrier.epoch;\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::Iceberg(iceberg_properties) = config else {\n            unreachable!()\n        };\n\n        let scan_projection =","sourceCodeStart":81,"sourceCodeEnd":117,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/source/iceberg_list_executor.rs#L81-L117","documentation":"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.","triggerScenarios":"into_stream() awaits barrier_receiver.recv() and gets None (channel closed / no senders) instead of the initial Barrier.","commonSituations":"Iceberg source actor cancelled at startup, compute node lost barrier connection to meta, or recovery/failover racing executor initialization.","solutions":["Check barrier manager connectivity from the compute node; recover the MV/job.","Inspect meta logs for cancellation or rescheduling of the source actor; retry creation.","In tests, send an initial Barrier into the channel before polling."],"exampleFix":null,"handlingStrategy":"retry","validationCode":"// preflight: ensure compute node has barrier manager connection before creating Iceberg sources","typeGuard":null,"tryCatchPattern":"match barrier_rx.recv().await {\n    Some(b) => Ok(b.epoch),\n    None => Err(recover_with_backoff(\"iceberg list: no first barrier\")),\n}","preventionTips":["Keep the cluster stable during CREATE SOURCE for Iceberg tables","Monitor barrier manager health and actor recovery events","In tests, send an initial barrier before polling the executor"],"tags":["streaming","barrier","iceberg"],"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"}