risingwavelabs/risingwave · error

sink fragment not found for sink

Error message

sink fragment not found for sink {}

What it means

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.

Solutions

  1. Recreate the sink so a fresh, well-formed plan is generated.
  2. Ensure frontend and meta versions match; redeploy consistent binaries.
  3. Inspect fragment metadata to see why no fragment qualifies as the sink fragment; fix classification or plan generation if in code.

Example fix

// before: sink plan split so sink_fragment() is None
// after
DROP SINK s;
CREATE SINK s INTO target_table AS SELECT * FROM mv; // fresh single sink fragment
Defensive patterns

Strategy: try-catch

Try / catch

match create_sink_into_table(..).await {
    Err(e) if e.to_string().contains("sink fragment not found") => {
        // drop and recreate the sink; check version alignment
    }
    other => other?,
}

Prevention

When it happens

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

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

Understand the failure class

Background: 'Could not be found', 'does not exist', 'not found in database': the resource-not-found family when an ID, slug, key, or URI lookup comes back empty — this error's family across 20 libraries.

Related errors


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

Appendix: source

Thrown at src/meta/src/rpc/ddl_controller.rs:2041

            StreamJobFragments::new(id, graph, stream_ctx.clone(), max_parallelism.get());

        if let Some(mview_fragment) = stream_job_fragments.mview_fragment() {
            stream_job.set_table_vnode_count(mview_fragment.vnode_count());
        }

        let new_upstream_sink = if let StreamingJob::Sink(sink, _) = &stream_job
            && let Ok(table_id) = sink.get_target_table()
        {
            let tables = self
                .metadata_manager
                .get_table_catalog_by_ids(&[*table_id])
                .await?;
            let target_table = tables
                .first()
                .ok_or_else(|| MetaError::catalog_id_not_found("table", *table_id))?;
            let sink_fragment = stream_job_fragments
                .sink_fragment()
                .ok_or_else(|| anyhow::anyhow!("sink fragment not found for sink {}", sink.id))?;
            let mview_fragment_id = self
                .metadata_manager
                .catalog_controller
                .get_mview_fragment_by_id(table_id.as_job_id())
                .await?;
            let upstream_sink_info = build_upstream_sink_info(
                sink.id,
                sink.original_target_columns.clone(),
                sink_fragment.fragment_id as _,
                target_table,
                mview_fragment_id,
            )?;
            Some(upstream_sink_info)
        } else {
            None
        };

        let mut cdc_table_snapshot_splits = None;

View on GitHub (pinned to 6469eb736d)