{"record":{"id":"01143e94a29c0ae0","repo":"risingwavelabs/risingwave","slug":"upstream-assignment-not-found-fragment-id-fragm","errorCode":null,"errorMessage":"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:?}","messagePattern":"upstream assignment not found, fragment_id: (.+?), upstream_fragment_id: (.+?), actor_id: (.+?), upstream_actor_id: (.+?)","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/meta/src/stream/source_manager/split_assignment.rs","lineNumber":627,"sourceCode":"/// illustration:\n/// ```text\n/// upstream                               new\n/// actor x1 [split 1, split2]      ->     actor y1 [split 1, split2]\n/// actor x2 [split 3]              ->     actor y2 [split 3]\n/// ...\n/// ```\npub fn align_splits(\n    // (actor_id, upstream_actor_id)\n    aligned_actors: impl IntoIterator<Item = (ActorId, ActorId)>,\n    get_upstream_actor_splits: impl Fn(ActorId) -> Option<Vec<SplitImpl>>,\n    fragment_id: FragmentId,\n    upstream_source_fragment_id: FragmentId,\n) -> anyhow::Result<HashMap<ActorId, Vec<SplitImpl>>> {\n    aligned_actors\n        .into_iter()\n        .map(|(actor_id, upstream_actor_id)| {\n            let Some(splits) = get_upstream_actor_splits(upstream_actor_id) else {\n                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:?}\"));\n            };\n\n            Ok((\n                actor_id,\n                splits,\n            ))\n        })\n        .collect()\n}\n\n/// Note: the `PartialEq` and `Ord` impl just compares the number of splits.\n#[derive(Debug)]\nstruct SplitsAssignment<I, T: SplitMetaData> {\n    actor_id: I,\n    splits: Vec<T>,\n}\n\nimpl<I, T: SplitMetaData + Clone> Eq for SplitsAssignment<I, T> {}","sourceCodeStart":609,"sourceCodeEnd":645,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/stream/source_manager/split_assignment.rs#L609-L645","documentation":"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.","triggerScenarios":"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).","commonSituations":"Source replacement or backfill actor migration racing with upstream actor creation; stale split assignment after meta restart; fragmented state after a failed rescale.","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"],"exampleFix":null,"handlingStrategy":"retry","validationCode":"// pre-check upstream assignments exist\nif get_upstream_actor_splits(upstream_actor_id).is_none() {\n    // wait for source manager tick before aligning splits\n}","typeGuard":"fn has_upstream_assignment(get: impl Fn(ActorId) -> Option<Vec<SplitImpl>>, id: ActorId) -> bool {\n    get(id).is_some()\n}","tryCatchPattern":"match align_splits(...) {\n    Err(e) if e.to_string().contains(\"upstream assignment not found\") => {\n        tokio::time::sleep(Duration::from_secs(5)).await; // retry after tick\n    }\n    r => r?,\n}","preventionTips":["Run split alignment only after upstream actors are started","Monitor source manager for missing assignments after meta restart"],"tags":["rust","meta","splits","source-manager"],"backgroundTag":"record-not-found","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"}