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
- 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.
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
- 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
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
- named already exists
- Dropping sink into table is not allowed for unmigrated table
- streaming jobs not found
- wait version not set
- id not found
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)