{"record":{"id":"2a77ea565d584ca9","repo":"risingwavelabs/risingwave","slug":"expected-exactly-one-sink-fragment-for-each-sink","errorCode":null,"errorMessage":"expected exactly one sink fragment for each sink, but got {} fragments for {} sinks","messagePattern":"expected exactly one sink fragment for each sink, but got (.+?) fragments for (.+?) sinks","errorType":"exception","errorClass":"MetaError","httpStatus":null,"severity":"error","filePath":"src/meta/src/controller/utils.rs","lineNumber":2389,"sourceCode":") -> MetaResult<HashMap<SinkId, FragmentId>>\nwhere\n    C: ConnectionTrait,\n{\n    let sink_num = sink_ids.len();\n    let sink_fragment_ids: Vec<(SinkId, FragmentId)> = Fragment::find()\n        .select_only()\n        .columns([fragment::Column::JobId, fragment::Column::FragmentId])\n        .filter(\n            fragment::Column::JobId\n                .is_in(sink_ids)\n                .and(FragmentTypeMask::intersects(FragmentTypeFlag::Sink)),\n        )\n        .into_tuple()\n        .all(txn)\n        .await?;\n\n    if sink_fragment_ids.len() != sink_num {\n        return Err(anyhow::anyhow!(\n            \"expected exactly one sink fragment for each sink, but got {} fragments for {} sinks\",\n            sink_fragment_ids.len(),\n            sink_num\n        )\n        .into());\n    }\n\n    Ok(sink_fragment_ids.into_iter().collect())\n}\n\npub async fn has_table_been_migrated<C>(txn: &C, table_id: TableId) -> MetaResult<bool>\nwhere\n    C: ConnectionTrait,\n{\n    let mview_fragment: Vec<i32> = Fragment::find()\n        .select_only()\n        .column(fragment::Column::FragmentTypeMask)\n        .filter(","sourceCodeStart":2371,"sourceCodeEnd":2407,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/controller/utils.rs#L2371-L2407","documentation":"This helper asserts a 1:1 mapping between sinks and their sink fragments in the catalog: the number of fetched sink fragment ids must equal the sink count. If they differ, it fails with \"expected exactly one sink fragment for each sink, but got {} fragments for {} sinks\", since downstream code assumes exactly one fragment per sink.","triggerScenarios":"Calling the sink-fragment lookup when some sinks have no persisted sink fragment row, or when duplicate/extra fragment rows exist — typically during schema changes, failover recovery, or a partially completed create/drop sink.","commonSituations":"Catalog corruption after an interrupted sink creation; version-upgrade migrations that changed fragment storage; internal tooling reading fragments mid-DDL; crash between catalog writes leaving sinks without fragments.","solutions":["Inspect the fragment/table catalogs for the affected sinks — delete orphan sinks or duplicate fragment rows, or recreate the sink cleanly.","Restore from backup if the catalog is inconsistent after a failed DDL; then re-create the sink.","Retry the operation after DDL settles; transient reads during sink creation/drop can see a mismatched count.","If this reproduces after an upgrade, run the upgrade/migration tooling that rebuilds fragment metadata."],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":"// Verify sink fragment rows before invoking the helper\nlet frags = fetch_sink_fragments(txn).await?;\nlet sinks = fetch_sink_count(txn).await?;\nif frags.len() != sinks { return Err(\"catalog inconsistent: sink fragment count mismatch\"); }","typeGuard":null,"tryCatchPattern":"match get_sink_fragments().await {\n    Err(e) if e.to_string().contains(\"one sink fragment for each sink\") => repair_or_recreate_sink(),\n    other => other,\n}","preventionTips":["Avoid reading fragment metadata while sink DDL is in flight","Clean up failed/interrupted sink creations promptly","After upgrades, verify catalog consistency with the provided migration checks"],"tags":["meta","sink","fragments","invariant"],"backgroundTag":"internal-invariant-violation","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"}