risingwavelabs/risingwave · error · DashboardError

Failed to get streaming stats from worker

Error message

Failed to get streaming stats from worker

What it means

The `get_streaming_stats_from_prometheus` dashboard handler also fans out GetStreamingStats RPCs to workers before consulting Prometheus. Any per-worker RPC failure is collapsed into this generic message; like its sibling path, the concrete gRPC error is discarded by `map_err(|_| ...)`.

Solutions

  1. Verify all workers are alive and connected (dashboard worker list / SHOW NODES) and restart unhealthy ones.
  2. Check meta node logs for the underlying per-worker error that this message hides.
  3. Retry after the cluster stabilizes.
  4. Patch the handler to include the original error in the message instead of `|_|`.

Example fix

// before
.map_err(|_| anyhow!("Failed to get streaming stats from worker"))
// after
.map_err(|e| anyhow!("Failed to get streaming stats from worker: {e}"))
Defensive patterns

Strategy: retry

Validate before calling

const health = await fetch('http://localhost:5691/api/v1/workers').then(r => r.json());
if (health.some(w => w.state !== 'RUNNING')) throw new Error('A worker is not RUNNING; streaming stats from prometheus path will fail');

Try / catch

try { const stats = await getStreamingStatsFromPrometheus(); } catch (e) {
  if (String(e).includes('Failed to get streaming stats from worker')) {
    // locate failing worker via meta logs, restart if needed, retry with backoff
  }
}

Prevention

When it happens

Trigger: GET /api/v1/streaming_stats (Prometheus-backed variant) while one or more workers fail the GetStreamingStats RPC — worker down, restarting, or unreachable from the meta node.

Common situations: Compute node OOM/killed during monitoring; network issues between meta and worker pods in Kubernetes; polling the dashboard mid-restart.

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/7f9ce5a9145e46c6. Report an issue: GitHub.

Appendix: source

Thrown at src/meta/src/dashboard/mod.rs:806

            .await
            .map_err(err)?;

        let mut futures = Vec::new();

        for worker_node in worker_nodes {
            let client = srv.monitor_clients.get(&worker_node).await.map_err(err)?;
            let client = Arc::new(client);
            let fut = async move {
                let result = client.get_streaming_stats().await.map_err(err)?;
                Ok::<_, DashboardError>(result)
            };
            futures.push(fut);
        }
        let results = join_all(futures).await;

        for result in results {
            let result = result
                .map_err(|_| anyhow!("Failed to get streaming stats from worker"))
                .map_err(err)?;

            // Aggregate fragment_stats
            for (fragment_id, fragment_stats) in result.fragment_stats {
                if let Some(s) = all.fragment_stats.get_mut(&fragment_id) {
                    s.actor_count += fragment_stats.actor_count;
                    s.current_epoch = min(s.current_epoch, fragment_stats.current_epoch);
                } else {
                    all.fragment_stats.insert(fragment_id, fragment_stats);
                }
            }

            // Aggregate relation_stats
            for (relation_id, relation_stats) in result.relation_stats {
                if let Some(s) = all.relation_stats.get_mut(&relation_id) {
                    s.actor_count += relation_stats.actor_count;
                    s.current_epoch = min(s.current_epoch, relation_stats.current_epoch);
                } else {

View on GitHub (pinned to 6469eb736d)