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
- 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.
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
- 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
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
- sink fragment not found for sink id
- ALTER SINK_RATE_LIMIT is not for sink into table
- ambiguous auth: multiple auth options provided; remove one…
- auth.method=key_pair_file must not set `password`
- auth.method=key_pair_file must not set `private_key_pem`
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)