{"record":{"id":"1af2d61cae269a34","repo":"risingwavelabs/risingwave","slug":"partial-graph-resetting","errorCode":null,"errorMessage":"partial graph resetting","messagePattern":"partial graph resetting","errorType":"exception","errorClass":"anyhow","httpStatus":null,"severity":"error","filePath":"src/stream/src/task/barrier_worker/mod.rs","lineNumber":552,"sourceCode":"                                );\n                                anyhow!(\n                                    \"take receiver {:?} on unmatched partial graph term {} to current term {}\",\n                                    ids,\n                                    term_id,\n                                    graph.local_barrier_manager.term_id\n                                )\n                            } else {\n                                let (upstream_actor_id, actor_id) = ids;\n                                graph.new_actor_output_request(\n                                    actor_id,\n                                    upstream_actor_id,\n                                    request,\n                                );\n                                return;\n                            }\n                        }\n                        PartialGraphStatus::Suspended(_) => anyhow!(\"partial graph suspended\"),\n                        PartialGraphStatus::Resetting => anyhow!(\"partial graph resetting\"),\n                        PartialGraphStatus::Unspecified => unreachable!(),\n                    },\n                    Entry::Vacant(entry) => {\n                        entry.insert(PartialGraphStatus::ReceivedExchangeRequest(vec![(\n                            term_id, ids, request,\n                        )]));\n                        return;\n                    }\n                };\n                if let TakeReceiverRequest::Remote { result_sender, .. } = request {\n                    let _ = result_sender.send(Err(err.into()));\n                }\n            }\n            #[cfg(test)]\n            LocalActorOperation::GetCurrentLocalBarrierManager(sender) => {\n                let partial_graph_status = self\n                    .state\n                    .partial_graphs","sourceCodeStart":534,"sourceCodeEnd":570,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/task/barrier_worker/mod.rs#L534-L570","documentation":"A TakeReceiver (exchange channel creation) request arrived while the target partial graph is being reset. Because the old graph is being torn down, the receiver cannot be taken and this error is returned to the remote requester. It is a transient, recovery-related rejection.","triggerScenarios":"handle_actor_op matches PartialGraphStatus::Resetting while processing LocalActorOperation::TakeReceiver for that partial graph id.","commonSituations":"Concurrent exchange channel setup during fault recovery; upstream retries racing with a ResetPartialGraphs command from the meta node.","solutions":["Retry the TakeReceiver request once the partial graph finishes resetting and is Running with a new term id.","Confirm the meta node issued ResetPartialGraphs and that the new term id is used in follow-up requests.","If this loops, check whether the reset itself keeps failing (look for ack_reset_partial_graph with a root error)."],"exampleFix":"// before: immediate take during reset\nworker.take_receiver(partial_graph_id, term_id, ids, request);\n// after: await ResetPartialGraph ack, then take with new term\nawait_reset_ack(partial_graph_id).await;\nworker.take_receiver(partial_graph_id, new_term_id, ids, request);","handlingStrategy":"retry","validationCode":"// avoid issuing take-receiver during reset window\nwhile matches!(graph_status(partial_graph_id), Resetting) { tokio::time::sleep(poll_interval).await; }","typeGuard":"fn is_resetting(err: &anyhow::Error) -> bool { err.to_string() == \"partial graph resetting\" }","tryCatchPattern":"// wait-and-retry on reset rejection\nmatch take_receiver().await {\n    Err(e) if is_resetting(&e) => { wait_reset_ack(partial_graph_id).await; take_receiver().await?; }\n    other => other,\n}","preventionTips":["Serialize graph creation/reset with exchange channel setup via the control stream ordering.","Treat Resetting as a normal transient state in retry logic, not a fatal error.","Monitor reset duration; long resets indicate recurring root failures."],"tags":["streaming","barrier","graph-reset","rpc"],"backgroundTag":"invalid-state-transition","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"}