{"record":{"id":"5824890b6452fbf5","repo":"risingwavelabs/risingwave","slug":"sink-fragment-not-found-for-sink","errorCode":null,"errorMessage":"sink fragment not found for sink {}","messagePattern":"sink fragment not found for sink (.+?)","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/meta/src/rpc/ddl_controller.rs","lineNumber":2041,"sourceCode":"            StreamJobFragments::new(id, graph, stream_ctx.clone(), max_parallelism.get());\n\n        if let Some(mview_fragment) = stream_job_fragments.mview_fragment() {\n            stream_job.set_table_vnode_count(mview_fragment.vnode_count());\n        }\n\n        let new_upstream_sink = if let StreamingJob::Sink(sink, _) = &stream_job\n            && let Ok(table_id) = sink.get_target_table()\n        {\n            let tables = self\n                .metadata_manager\n                .get_table_catalog_by_ids(&[*table_id])\n                .await?;\n            let target_table = tables\n                .first()\n                .ok_or_else(|| MetaError::catalog_id_not_found(\"table\", *table_id))?;\n            let sink_fragment = stream_job_fragments\n                .sink_fragment()\n                .ok_or_else(|| anyhow::anyhow!(\"sink fragment not found for sink {}\", sink.id))?;\n            let mview_fragment_id = self\n                .metadata_manager\n                .catalog_controller\n                .get_mview_fragment_by_id(table_id.as_job_id())\n                .await?;\n            let upstream_sink_info = build_upstream_sink_info(\n                sink.id,\n                sink.original_target_columns.clone(),\n                sink_fragment.fragment_id as _,\n                target_table,\n                mview_fragment_id,\n            )?;\n            Some(upstream_sink_info)\n        } else {\n            None\n        };\n\n        let mut cdc_table_snapshot_splits = None;","sourceCodeStart":2023,"sourceCodeEnd":2059,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/rpc/ddl_controller.rs#L2023-L2059","documentation":"When building a sink-into-table streaming job, build_stream_job needs the sink's own fragment to construct upstream sink info. If stream_job_fragments.sink_fragment() returns None, the sink plan has no identifiable sink fragment and the build fails, naming the sink ID.","triggerScenarios":"Creating/updating a sink into a table where the generated fragment graph has no fragment classified as the sink fragment — e.g. plan shape changes, wrong sink topology, or fragment mislabeling.","commonSituations":"Sink plans altered by frontend changes so the sink operator no longer lands in the expected fragment; version skew between frontend and meta; corrupted or partially-created sink jobs being rebuilt.","solutions":["Recreate the sink so a fresh, well-formed plan is generated.","Ensure frontend and meta versions match; redeploy consistent binaries.","Inspect fragment metadata to see why no fragment qualifies as the sink fragment; fix classification or plan generation if in code."],"exampleFix":"// before: sink plan split so sink_fragment() is None\n// after\nDROP SINK s;\nCREATE SINK s INTO target_table AS SELECT * FROM mv; // fresh single sink fragment","handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"match create_sink_into_table(..).await {\n    Err(e) if e.to_string().contains(\"sink fragment not found\") => {\n        // drop and recreate the sink; check version alignment\n    }\n    other => other?,\n}","preventionTips":["Keep frontend and meta versions consistent","Recreate sinks whose plans may predate current plan shapes","Inspect fragment view to confirm sink fragment classification before rebuilds"],"tags":["sink","fragment","plan"],"backgroundTag":"resource-not-found","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"}