risingwavelabs/risingwave · error · MetaError
anyhow!(message.to_owned())
Error message
anyhow!(message.to_owned())
What it means
merge_node_rpc_errors collapses a set of per-worker-node RPC errors into a single MetaError. When the caller passes an empty error list, the function has nothing to attribute to any worker, so it throws a bare anyhow error containing only the generic message. It exists so the meta service can return one consolidated barrier/RPC failure instead of leaking many per-node errors to the frontend.
Source
Thrown at src/meta/src/barrier/rpc.rs:1495
.start_streaming_control(PbInitRequest::default())
.await?;
Ok(handle)
}
}
pub(super) fn merge_node_rpc_errors<E: Error + Send + Sync + 'static>(
message: &str,
errors: impl IntoIterator<Item = (WorkerId, E)>,
) -> MetaError {
use std::fmt::Write;
use risingwave_common::error::error_request_copy;
use risingwave_common::error::tonic::extra::Score;
let errors = errors.into_iter().collect_vec();
if errors.is_empty() {
return anyhow!(message.to_owned()).into();
}
// Create the error from the single error.
let single_error = |(worker_id, e)| {
anyhow::Error::from(e)
.context(format!("{message}, in worker node {worker_id}"))
.into()
};
if errors.len() == 1 {
return single_error(errors.into_iter().next().unwrap());
}
// Find the error with the highest score.
let max_score = errors
.iter()
.filter_map(|(_, e)| error_request_copy::<Score>(e))
.max();View on GitHub (pinned to 6469eb736d)
Solutions
- Check cluster health: ensure streaming worker nodes are registered and alive before triggering barrier operations (`SHOW CLUSTERS` / meta node logs for worker liveness).
- If operating during scale-down, wait for or trigger recovery so the barrier manager re-establishes its worker set.
- Upgrade/Rollback recently: if this appeared after a version change, restart the cluster with recovery enabled so worker-node state is rebuilt.
Example fix
// before: merging with an empty error list
let errors: Vec<_> = collected_errors.into_iter().filter(|e| !is_benign(e)).collect();
let merged = merge_node_rpc_errors(msg, errors);
// after: skip merging entirely when nothing failed
let errors: Vec<_> = collected_errors.into_iter().filter(|e| !is_benign(e)).collect();
let merged = if errors.is_empty() { Ok(()) } else { merge_node_rpc_errors(msg, errors) }; Defensive patterns
Strategy: retry
Try / catch
// Rust
match result {
Err(e) if is_empty_aggregate_rpc_error(&e) => warn!("no workers targeted; retry after recovery"),
Err(e) => return Err(e),
Ok(v) => Ok(v),
} Prevention
- Monitor worker liveness and never let the cluster run with zero streaming workers.
- Keep recovery enabled so worker loss triggers automatic re-registration.
- Check `SHOW CLUSTERS;` before maintenance that could drain all compute nodes.
When it happens
Trigger: Calling merge_node_rpc_errors (barrier dispatch paths in the meta node) with an empty `errors` iterator — i.e. the fan-out to worker nodes produced zero errors yet the caller still took the error-merging branch, typically when no worker nodes were targeted or all error entries were filtered out before merging.
Common situations: Cluster where all streaming workers have been removed/failed but the barrier manager still attempts to dispatch; races between worker deregistration and barrier issuance; tests or internal callers invoking the merge helper without collected errors.
Related errors
- {message}: in worker node {}, {};
- since_timestamp epoch has not been resolved for snapshot bac
- cannot create batch refresh job while database barrier is pa
- replace sink must not use snapshot backfill
- old sink job {} not found in barrier state
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/d27379b77274050e.
Report an issue: GitHub.