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
- Retry the source operation after upstream actors are running
- Verify upstream actor split assignments exist (source manager state/logs)
- Trigger a split reassignment/tick so the source manager repopulates assignments
- 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
- Run split alignment only after upstream actors are started
- Monitor source manager for missing assignments after meta restart
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
- source not found in source manager
- cannot find StreamActor of actor
- concurrent backup job is not supported: existent job
- `debug_splits` is not allowed in release mode
- downstream relation missing for
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)