{"record":{"id":"97956a9c13d804a2","repo":"risingwavelabs/risingwave","slug":"replace-sink-requires-a-sink-job","errorCode":null,"errorMessage":"replace sink requires a sink job","messagePattern":"replace sink requires a sink job","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/meta/src/rpc/ddl_controller.rs","lineNumber":1115,"sourceCode":"    }\n\n    /// For [`CreateType::Foreground`], the function will only return after backfilling finishes\n    /// ([`crate::manager::MetadataManager::wait_streaming_job_finished`]).\n    #[await_tree::instrument(boxed, \"create_streaming_job({streaming_job})\")]\n    pub async fn create_streaming_job(\n        &self,\n        mut streaming_job: StreamingJob,\n        fragment_graph: StreamFragmentGraphProto,\n        dependencies: HashSet<ObjectId>,\n        resource_type: streaming_job_resource_type::ResourceType,\n        if_not_exists: bool,\n        refresh_interval_sec: Option<u64>,\n        replace_sink: Option<SinkId>,\n        since_timestamp_epoch: Option<u64>,\n    ) -> MetaResult<NotificationVersion> {\n        let replace_sink_info = if let Some(old_sink_id) = replace_sink {\n            let StreamingJob::Sink(sink, _) = &streaming_job else {\n                bail!(\"replace sink requires a sink job\")\n            };\n            if sink.target_table.is_some() {\n                bail_not_implemented!(\"replace sink into table\")\n            }\n\n            Some(old_sink_id)\n        } else {\n            if let StreamingJob::Sink(sink, _) = &streaming_job\n                && let Some(target_table) = sink.target_table\n            {\n                self.validate_table_for_sink(target_table).await?;\n            }\n            None\n        };\n        self.validate_serverless_backfill_enabled(&resource_type)?;\n        let ctx = StreamContext::from_protobuf(fragment_graph.get_ctx().unwrap());\n        let adaptive_parallelism_strategy =\n            (!fragment_graph.adaptive_parallelism_strategy.is_empty()).then(|| {","sourceCodeStart":1097,"sourceCodeEnd":1133,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/rpc/ddl_controller.rs#L1097-L1133","documentation":"The replace_sink path updates an existing sink by resubmitting a streaming job, but that job must actually be a sink. create_streaming_job checks that when replace_sink is Some, the streaming_job matches StreamingJob::Sink; otherwise the operation is invalid and bails.","triggerScenarios":"Calling the DDL controller's create_streaming_job with a non-sink StreamingJob (table, index, MV) while passing replace_sink = Some(old_sink_id).","commonSituations":"Client/driver bugs mapping ALTER SINK to the wrong streaming job type; custom tooling reusing the create API for replacement; refactors that change StreamingJob variants without updating replace logic.","solutions":["Pass a StreamingJob::Sink when using replace_sink; use a plain create path for non-sink jobs.","Fix the calling code so ALTER SINK builds the correct StreamingJob::Sink variant.","If you meant to replace a table/index job, use the dedicated ALTER/replace paths for those objects."],"exampleFix":"// before\ncontroller.create_streaming_job(StreamingJob::Table(..), None, Some(sink_id), ..).await?;\n// after\ncontroller.create_streaming_job(StreamingJob::Sink(sink, ..), None, Some(sink_id), ..).await?;","handlingStrategy":"validation","validationCode":"// caller-side check\nif replace_sink.is_some() && !matches!(streaming_job, StreamingJob::Sink(..)) {\n    return Err(\"replace_sink requires StreamingJob::Sink\".into());\n}","typeGuard":"fn is_sink_job(job: &StreamingJob) -> bool {\n    matches!(job, StreamingJob::Sink(..))\n}","tryCatchPattern":"match create_streaming_job(job, ..).await {\n    Err(e) if e.to_string().contains(\"replace sink requires a sink job\") => fix_job_type(),\n    other => other?,\n}","preventionTips":["Only pass replace_sink when the job is genuinely a sink","Type-check StreamingJob variant before invoking replace paths","Use dedicated replace/alter APIs for non-sink objects"],"tags":["sink","api-misuse","validation"],"backgroundTag":"invalid-argument-value","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"}