risingwavelabs/risingwave · error · MetaError

expected exactly one sink fragment for each sink, but got

Error message

expected exactly one sink fragment for each sink, but got {} fragments for {} sinks

What it means

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.

Solutions

  1. Inspect the fragment/table catalogs for the affected sinks — delete orphan sinks or duplicate fragment rows, or recreate the sink cleanly.
  2. Restore from backup if the catalog is inconsistent after a failed DDL; then re-create the sink.
  3. Retry the operation after DDL settles; transient reads during sink creation/drop can see a mismatched count.
  4. If this reproduces after an upgrade, run the upgrade/migration tooling that rebuilds fragment metadata.
Defensive patterns

Strategy: try-catch

Validate before calling

// Verify sink fragment rows before invoking the helper
let frags = fetch_sink_fragments(txn).await?;
let sinks = fetch_sink_count(txn).await?;
if frags.len() != sinks { return Err("catalog inconsistent: sink fragment count mismatch"); }

Try / catch

match get_sink_fragments().await {
    Err(e) if e.to_string().contains("one sink fragment for each sink") => repair_or_recreate_sink(),
    other => other,
}

Prevention

When it happens

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

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

Understand the failure class

Background: "This is a bug, please report it": internal invariant violations, unreachable panics, and SNH errors explained — this error's family across 47 libraries.

Related errors


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

Appendix: source

Thrown at src/meta/src/controller/utils.rs:2389

) -> MetaResult<HashMap<SinkId, FragmentId>>
where
    C: ConnectionTrait,
{
    let sink_num = sink_ids.len();
    let sink_fragment_ids: Vec<(SinkId, FragmentId)> = Fragment::find()
        .select_only()
        .columns([fragment::Column::JobId, fragment::Column::FragmentId])
        .filter(
            fragment::Column::JobId
                .is_in(sink_ids)
                .and(FragmentTypeMask::intersects(FragmentTypeFlag::Sink)),
        )
        .into_tuple()
        .all(txn)
        .await?;

    if sink_fragment_ids.len() != sink_num {
        return Err(anyhow::anyhow!(
            "expected exactly one sink fragment for each sink, but got {} fragments for {} sinks",
            sink_fragment_ids.len(),
            sink_num
        )
        .into());
    }

    Ok(sink_fragment_ids.into_iter().collect())
}

pub async fn has_table_been_migrated<C>(txn: &C, table_id: TableId) -> MetaResult<bool>
where
    C: ConnectionTrait,
{
    let mview_fragment: Vec<i32> = Fragment::find()
        .select_only()
        .column(fragment::Column::FragmentTypeMask)
        .filter(

View on GitHub (pinned to 6469eb736d)