risingwavelabs/risingwave · error

failed to wait streaming job finish

Error message

failed to wait streaming job finish: {}

What it means

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)`.

Solutions

  1. Inspect the embedded `reason` string in the message — it names the actual failure of the streaming job and fix that root cause.
  2. Retry the DDL statement (CREATE MV/SINK) after resolving the underlying job failure.
  3. Check meta-node logs for the job_id to see whether the job crashed, was cancelled, or the notifier was dropped.
  4. Ensure meta-node is not restarted/crashed mid-DDL; verify cluster stability before long-running streaming jobs.
Defensive patterns

Strategy: try-catch

Validate before calling

// check the job is still tracked before waiting
if !meta.list_creating_jobs().await?.contains(&job_id) { /* job already finished or lost */ }

Try / catch

match meta.wait_streaming_job_finished(db_id, job_id).await {
    Ok(_) => (),
    Err(e) => {
        let reason = e.to_string();
        tracing::error!("streaming job {} failed: {}", job_id, reason);
        // inspect `reason` for the underlying job failure and retry DDL if transient
    }
}

Prevention

When it happens

Trigger: 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.

Common situations: 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.

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/a91e3b5cb3df1a02. Report an issue: GitHub.

Appendix: source

Thrown at src/meta/src/manager/metadata.rs:812

    #[await_tree::instrument]
    pub async fn wait_streaming_job_finished(
        &self,
        database_id: DatabaseId,
        id: JobId,
    ) -> MetaResult<NotificationVersion> {
        tracing::debug!("wait_streaming_job_finished: {id:?}");
        let mut mgr = self.catalog_controller.get_inner_write_guard().await;
        if mgr.streaming_job_is_finished(id).await? {
            return Ok(self.catalog_controller.notify_frontend_trivial().await);
        }
        let (tx, rx) = oneshot::channel();

        mgr.register_finish_notifier(database_id, id, tx);
        drop(mgr);
        rx.await
            .map_err(|_| "no received reason".to_owned())
            .and_then(|result| result)
            .map_err(|reason| anyhow!("failed to wait streaming job finish: {}", reason).into())
    }

    pub(crate) async fn notify_finish_failed(&self, database_id: Option<DatabaseId>, err: String) {
        let mut mgr = self.catalog_controller.get_inner_write_guard().await;
        mgr.notify_finish_failed(database_id, err);
    }
}

View on GitHub (pinned to 6469eb736d)