{"record":{"id":"a91e3b5cb3df1a02","repo":"risingwavelabs/risingwave","slug":"failed-to-wait-streaming-job-finish","errorCode":null,"errorMessage":"failed to wait streaming job finish: {}","messagePattern":"failed to wait streaming job finish: (.+?)","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/meta/src/manager/metadata.rs","lineNumber":812,"sourceCode":"    #[await_tree::instrument]\n    pub async fn wait_streaming_job_finished(\n        &self,\n        database_id: DatabaseId,\n        id: JobId,\n    ) -> MetaResult<NotificationVersion> {\n        tracing::debug!(\"wait_streaming_job_finished: {id:?}\");\n        let mut mgr = self.catalog_controller.get_inner_write_guard().await;\n        if mgr.streaming_job_is_finished(id).await? {\n            return Ok(self.catalog_controller.notify_frontend_trivial().await);\n        }\n        let (tx, rx) = oneshot::channel();\n\n        mgr.register_finish_notifier(database_id, id, tx);\n        drop(mgr);\n        rx.await\n            .map_err(|_| \"no received reason\".to_owned())\n            .and_then(|result| result)\n            .map_err(|reason| anyhow!(\"failed to wait streaming job finish: {}\", reason).into())\n    }\n\n    pub(crate) async fn notify_finish_failed(&self, database_id: Option<DatabaseId>, err: String) {\n        let mut mgr = self.catalog_controller.get_inner_write_guard().await;\n        mgr.notify_finish_failed(database_id, err);\n    }\n}\n","sourceCodeStart":794,"sourceCodeEnd":820,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/manager/metadata.rs#L794-L820","documentation":"Returned by `wait_streaming_job_finished` when the oneshot channel that reports job completion fails: either the sender was dropped without sending ('no received reason') or the finishing path delivered an error reason string, which is wrapped via `anyhow!(\"failed to wait streaming job finish: {}\", reason)`.","triggerScenarios":"Calling `wait_streaming_job_finished` and awaiting `rx` when the catalog manager dropped the registered finish notifier without notifying (e.g. job cleanup, meta-node shutdown), or when `notify_finish_failed` / a failed finish path sends an error reason.","commonSituations":"Streaming job creation fails during DDL execution; meta-node restarts while a CREATE MATERIALIZED VIEW / sink is still in-flight; the finish notification path is skipped due to a bug or race in job finalization.","solutions":["Inspect the embedded `reason` string in the message — it names the actual failure of the streaming job and fix that root cause.","Retry the DDL statement (CREATE MV/SINK) after resolving the underlying job failure.","Check meta-node logs for the job_id to see whether the job crashed, was cancelled, or the notifier was dropped.","Ensure meta-node is not restarted/crashed mid-DDL; verify cluster stability before long-running streaming jobs."],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":"// check the job is still tracked before waiting\nif !meta.list_creating_jobs().await?.contains(&job_id) { /* job already finished or lost */ }","typeGuard":null,"tryCatchPattern":"match meta.wait_streaming_job_finished(db_id, job_id).await {\n    Ok(_) => (),\n    Err(e) => {\n        let reason = e.to_string();\n        tracing::error!(\"streaming job {} failed: {}\", job_id, reason);\n        // inspect `reason` for the underlying job failure and retry DDL if transient\n    }\n}","preventionTips":["Always log the embedded reason string — it names the real job failure","Avoid restarting meta nodes while DDL jobs are in flight","Monitor creating-jobs and finish notifiers to detect dropped notifiers early"],"tags":["meta","streaming-job","ddl"],"backgroundTag":"unexpected-response-shape","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"}