risingwavelabs/risingwave · warning · anyhow
take receiver ( , ) on stale partial graph term ; current…
Error message
take receiver ({upstream_actor_id}, {actor_id}) on stale partial graph term {request_term_id}; current term is {term_id} What it means
When a new partial graph is created (add_partial_graph), the worker replays buffered exchange TakeReceiver requests. Any buffered request whose term id is older than the newly created graph's term id is stale: it is rejected by sending this error to the remote requester's oneshot result_sender. This prevents wiring channels from an obsolete graph generation into the fresh one.
Solutions
- Have the upstream/downstream retry the channel creation with the current term id after CreatePartialGraph.
- Verify term id propagation from the meta node so actors always request with the latest term.
- Check for delayed or replayed control requests (network retries, stale gRPC buffers) that carry old terms.
- Treat this rejection as expected during recovery if followed by a successful retry.
Example fix
// before: replaying all buffered requests verbatim
for (term_id, ids, request) in pending { dispatch(request); }
// after: only replay requests whose term matches the new graph
for (req_term, ids, request) in pending {
if req_term == term_id { dispatch(request); } else { reject_stale(request, req_term, term_id); }
} Defensive patterns
Strategy: retry
Validate before calling
// drop buffered requests whose term predates the new graph let live: Vec<_> = pending.into_iter().filter(|(t, _, _)| *t == term_id).collect();
Type guard
fn is_stale_term(err: &str) -> bool { err.contains("stale partial graph term") } Try / catch
// remote side: on stale rejection, re-request with current term
if result.is_err() && is_stale_term(&result.err().unwrap().to_string()) {
reissue_take_receiver_with_latest_term().await;
} Prevention
- Discard in-flight exchange requests whenever a term bump is observed.
- Bound retry counts for stale-term requests to avoid unbounded loops.
- Keep term id monotonic and propagated through every control message.
When it happens
Trigger: add_partial_graph (called via handle_streaming_control_request for Request::CreatePartialGraph) drains buffered ReceivedExchangeRequest entries and finds request_term_id != term_id of the new graph, rejecting TakeReceiverRequest::Remote via result_sender.
Common situations: Upstream exchange requests buffered while the graph was absent, then superseded by a graph rebuild with a bumped term; slow upstream actor retrying an epoch that was already replaced during recovery.
Understand the failure class
Background: "Invalid state transition" errors: "status must be X, actually Y", "already rejected/charging/uninstalled", "cannot ... while running" — what they mean when a library rejects your call — this error's family across 31 libraries.
Related errors
- take receiver on unmatched partial graph term to current…
- actor exited unexpectedly
- barrier reader closed unexpectedly
- current epoch has exceeded the epoch of the stream that has…
- end of barrier receiver
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/2c194ad6ea264a32.
Report an issue: GitHub.
Appendix: source
Thrown at src/stream/src/task/barrier_worker/mod.rs:965
term_id.clone(),
self.actor_manager.clone(),
);
for (request_term_id, (upstream_actor_id, actor_id), request) in
pending_requests.drain(..)
{
if request_term_id == term_id {
graph.new_actor_output_request(actor_id, upstream_actor_id, request);
} else {
warn!(
%partial_graph_id,
%upstream_actor_id,
%actor_id,
request_term_id,
current_term_id = term_id,
"reject buffered exchange request with stale partial graph term"
);
if let TakeReceiverRequest::Remote { result_sender, .. } = request {
let _ = result_sender.send(Err(anyhow!(
"take receiver ({upstream_actor_id}, {actor_id}) on stale partial graph term {request_term_id}; current term is {term_id}"
)
.into()));
}
}
}
*status = PartialGraphStatus::Running(graph);
} else {
panic!("duplicated partial graph: {}", partial_graph_id);
}
status
}
Entry::Vacant(entry) => entry.insert(PartialGraphStatus::Running(
PartialGraphState::new(partial_graph_id, term_id, self.actor_manager.clone()),
)),
};
}View on GitHub (pinned to 6469eb736d)