{"record":{"id":"fcef6a0294469d05","repo":"risingwavelabs/risingwave","slug":"failed-to-send-shutdown-signal-for-streaming-job","errorCode":null,"errorMessage":"failed to send shutdown signal for streaming job {}: receiver dropped","messagePattern":"failed to send shutdown signal for streaming job (.+?): receiver dropped","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/meta/src/stream/stream_manager.rs","lineNumber":222,"sourceCode":"    }\n\n    async fn cancel_jobs(\n        &self,\n        job_ids: Vec<JobId>,\n    ) -> MetaResult<(HashMap<JobId, oneshot::Receiver<bool>>, Vec<JobId>)> {\n        let mut jobs = self.streaming_jobs.lock().await;\n        let mut receivers = HashMap::new();\n        let mut background_job_ids = vec![];\n        for job_id in job_ids {\n            if let Some(job) = jobs.get_mut(&job_id) {\n                if let Some(shutdown_tx) = job.shutdown_tx.take() {\n                    let (tx, rx) = oneshot::channel();\n                    match shutdown_tx.send(tx) {\n                        Ok(()) => {\n                            receivers.insert(job_id, rx);\n                        }\n                        Err(_) => {\n                            return Err(anyhow::anyhow!(\n                                \"failed to send shutdown signal for streaming job {}: receiver dropped\",\n                                job_id\n                            )\n                            .into());\n                        }\n                    }\n                }\n            } else {\n                // If these job ids do not exist in streaming_jobs, they should be background creating jobs.\n                background_job_ids.push(job_id);\n            }\n        }\n\n        Ok((receivers, background_job_ids))\n    }\n}\n\ntype CreatingStreamingJobInfoRef = Arc<CreatingStreamingJobInfo>;","sourceCodeStart":204,"sourceCodeEnd":240,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/stream/stream_manager.rs#L204-L240","documentation":"In `cancel_jobs`, the meta node sends a shutdown oneshot sender (`tx`) over the job shutdown channel to the component that owns streaming job lifecycle. If that channel's receiver has been dropped, the send fails and canceling the job returns this error, meaning the shutdown request could not be delivered.","triggerScenarios":"Calling cancel/drop for a streaming job while the shutdown signal receiver (held by the stream manager/actor shutdown task) is no longer alive — e.g. the shutdown worker already exited, meta node is shutting down, or the receiver task crashed.","commonSituations":"Dropping a MV/table while the meta node is concurrently shutting down or restarting; internal crashes of the shutdown handling task; heavy load causing the receiver task to be cancelled.","solutions":["Retry the cancellation once the meta node is healthy; verify with SHOW MATERIALIZED VIEWS / catalog whether the job was actually dropped.","Check meta node logs for shutdown task crashes and restart the meta service if the shutdown channel is permanently broken.","If the job is already gone but the error surfaces, treat it as benign; ensure cleanup paths tolerate `send` failure instead of failing the cancel."],"exampleFix":"// before\nErr(_) => return Err(anyhow::anyhow!(\"failed to send shutdown signal for streaming job {}: receiver dropped\", job_id).into()),\n// after: treat dropped receiver as job already shut down\nErr(_) => {\n    tracing::warn!(\"shutdown receiver for job {} dropped; assuming already stopped\", job_id);\n}","handlingStrategy":"retry","validationCode":"// Check the shutdown channel is alive before cancelling\nif shutdown_rx.is_closed() {\n    return Err(\"shutdown receiver unavailable; meta may be shutting down\");\n}","typeGuard":null,"tryCatchPattern":"match cancel_streaming_job(job_id).await {\n    Err(e) if e.to_string().contains(\"receiver dropped\") => {\n        // receiver gone: verify job state, then retry after meta recovers\n        if job_still_exists(job_id).await { retry_cancel(job_id).await?; }\n    }\n    other => other?,\n}","preventionTips":["Avoid cancelling streaming jobs during meta node shutdown/restart","Monitor meta logs for shutdown task crashes","Make cancellation idempotent: check job existence before and after cancel"],"tags":["meta","streaming","shutdown","lifecycle"],"backgroundTag":"broken-pipe","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"}