risingwavelabs/risingwave · error

upstream assignment not found, fragment_id

Error message

upstream assignment not found, fragment_id: {fragment_id}, upstream_fragment_id: {upstream_source_fragment_id}, actor_id: {actor_id}, upstream_actor_id: {upstream_actor_id:?}

What it means

align_splits maps each downstream actor to its upstream actor and then looks up that upstream actor's current split assignment via get_upstream_actor_splits; this error fires when the upstream actor has no recorded split assignment. It propagates through reassign_splits, migrate_splits_for_backfill_actors, resolve_replace_source_splits, and resolve_backfill_splits.

Solutions

  1. Retry the source operation after upstream actors are running
  2. Verify upstream actor split assignments exist (source manager state/logs)
  3. Trigger a split reassignment/tick so the source manager repopulates assignments
  4. Restart the meta source manager or recreate the source if state is corrupt
Defensive patterns

Strategy: retry

Validate before calling

// pre-check upstream assignments exist
if get_upstream_actor_splits(upstream_actor_id).is_none() {
    // wait for source manager tick before aligning splits
}

Type guard

fn has_upstream_assignment(get: impl Fn(ActorId) -> Option<Vec<SplitImpl>>, id: ActorId) -> bool {
    get(id).is_some()
}

Try / catch

match align_splits(...) {
    Err(e) if e.to_string().contains("upstream assignment not found") => {
        tokio::time::sleep(Duration::from_secs(5)).await; // retry after tick
    }
    r => r?,
}

Prevention

When it happens

Trigger: Split reassignment/alignment when an upstream actor referenced by the no_shuffle mapping has no entry in the upstream assignment map (e.g. upstream actor not yet started, or assignment dropped).

Common situations: Source replacement or backfill actor migration racing with upstream actor creation; stale split assignment after meta restart; fragmented state after a failed rescale.

Understand the failure class

Background: Record Not Found Errors: "not found", RecordNotFound, and "was not found" — what they mean and how to fix them — this error's family across 28 libraries.

Related errors


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

Appendix: source

Thrown at src/meta/src/stream/source_manager/split_assignment.rs:627

/// illustration:
/// ```text
/// upstream                               new
/// actor x1 [split 1, split2]      ->     actor y1 [split 1, split2]
/// actor x2 [split 3]              ->     actor y2 [split 3]
/// ...
/// ```
pub fn align_splits(
    // (actor_id, upstream_actor_id)
    aligned_actors: impl IntoIterator<Item = (ActorId, ActorId)>,
    get_upstream_actor_splits: impl Fn(ActorId) -> Option<Vec<SplitImpl>>,
    fragment_id: FragmentId,
    upstream_source_fragment_id: FragmentId,
) -> anyhow::Result<HashMap<ActorId, Vec<SplitImpl>>> {
    aligned_actors
        .into_iter()
        .map(|(actor_id, upstream_actor_id)| {
            let Some(splits) = get_upstream_actor_splits(upstream_actor_id) else {
                return Err(anyhow::anyhow!("upstream assignment not found, fragment_id: {fragment_id}, upstream_fragment_id: {upstream_source_fragment_id}, actor_id: {actor_id}, upstream_actor_id: {upstream_actor_id:?}"));
            };

            Ok((
                actor_id,
                splits,
            ))
        })
        .collect()
}

/// Note: the `PartialEq` and `Ord` impl just compares the number of splits.
#[derive(Debug)]
struct SplitsAssignment<I, T: SplitMetaData> {
    actor_id: I,
    splits: Vec<T>,
}

impl<I, T: SplitMetaData + Clone> Eq for SplitsAssignment<I, T> {}

View on GitHub (pinned to 6469eb736d)