risingwavelabs/risingwave · error

cannot get database_id of fragment

Error message

cannot get database_id of fragment {fragment_id}

What it means

Thrown by `split_fragment_map_by_database` when a fragment id present in the fragment map has no corresponding entry in the `fragment_to_database_map` derived from catalog tables. This is a catalog-consistency check: every streaming fragment must belong to a known database. It surfaces as an anyhow::Error propagated to the caller of the public API.

Solutions

  1. Check catalog consistency: verify the fragment exists in the fragment catalog table together with its database_id for the same snapshot/version.
  2. Rebuild `fragment_to_database_map` and the fragment map from the same catalog snapshot so they cannot diverge.
  3. If stale fragments from failed jobs exist, clean them up (or skip unmapped fragments) before splitting.
  4. Upgrade/repair meta-node state if the catalog is inconsistent after failover or restore.

Example fix

// before: two independently fetched maps may diverge
let fragment_to_database_map = build_fragment_to_database_map().await;
let fragment_map = list_fragments().await;
// after: fetch both from one consistent snapshot
let snapshot = catalog.snapshot().await;
let fragment_to_database_map = build_fragment_to_database_map(&snapshot);
let fragment_map = list_fragments(&snapshot);
Defensive patterns

Strategy: validation

Validate before calling

// ensure every fragment id is mapped before splitting
let unmapped: Vec<_> = fragment_map.keys().filter(|id| !fragment_to_database_map.contains_key(*id)).collect();
if !unmapped.is_empty() { return Err(anyhow!("unmapped fragments: {:?}", unmapped)); }

Type guard

fn is_mapped(fragment_id: &FragmentId, map: &HashMap<u32, DatabaseId>) -> bool { map.contains_key(fragment_id) }

Try / catch

match split_fragment_map_by_database(...).await {
    Ok(map) => use(map),
    Err(e) if e.to_string().contains("cannot get database_id") => refresh_catalog_snapshot_and_retry(),
    Err(e) => return Err(e),
}

Prevention

When it happens

Trigger: Calling `split_fragment_map_by_database` with a fragment map containing a fragment_id that is absent from `fragment_to_database_map` — e.g. the fragment was created but its catalog database row was not yet visible, or the two maps were built from different snapshot points.

Common situations: Concurrent DDL while a meta-node snapshot is taken; stale or partially-loaded catalog after failover; fragment records left behind by a failed streaming job whose database mapping was cleaned up; bugs in catalog version synchronization between readers.

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/38b52ff13b6ec322. Report an issue: GitHub.

Appendix: source

Thrown at src/meta/src/manager/metadata.rs:348

        fragment_map: HashMap<FragmentId, T>,
    ) -> MetaResult<HashMap<DatabaseId, HashMap<FragmentId, T>>> {
        let fragment_to_database_map: HashMap<_, _> = self
            .catalog_controller
            .list_fragment_database_ids(Some(
                fragment_map
                    .keys()
                    .map(|fragment_id| *fragment_id as _)
                    .collect(),
            ))
            .await?
            .into_iter()
            .map(|(fragment_id, database_id)| (fragment_id as FragmentId, database_id))
            .collect();
        let mut ret: HashMap<_, HashMap<_, _>> = HashMap::new();
        for (fragment_id, value) in fragment_map {
            let database_id = *fragment_to_database_map
                .get(&fragment_id)
                .ok_or_else(|| anyhow!("cannot get database_id of fragment {fragment_id}"))?;
            ret.entry(database_id)
                .or_default()
                .try_insert(fragment_id, value)
                .expect("non duplicate");
        }
        Ok(ret)
    }

    pub async fn list_creating_jobs(&self) -> MetaResult<HashSet<JobId>> {
        Ok(self
            .catalog_controller
            .list_creating_jobs(false, None)
            .await?
            .into_iter()
            .map(|(job_id, _, _, _, _)| job_id)
            .collect())
    }

View on GitHub (pinned to 6469eb736d)