{"record":{"id":"7f9ce5a9145e46c6","repo":"risingwavelabs/risingwave","slug":"failed-to-get-streaming-stats-from-worker","errorCode":null,"errorMessage":"Failed to get streaming stats from worker","messagePattern":"Failed to get streaming stats from worker","errorType":"http","errorClass":"DashboardError","httpStatus":500,"severity":"error","filePath":"src/meta/src/dashboard/mod.rs","lineNumber":806,"sourceCode":"            .await\n            .map_err(err)?;\n\n        let mut futures = Vec::new();\n\n        for worker_node in worker_nodes {\n            let client = srv.monitor_clients.get(&worker_node).await.map_err(err)?;\n            let client = Arc::new(client);\n            let fut = async move {\n                let result = client.get_streaming_stats().await.map_err(err)?;\n                Ok::<_, DashboardError>(result)\n            };\n            futures.push(fut);\n        }\n        let results = join_all(futures).await;\n\n        for result in results {\n            let result = result\n                .map_err(|_| anyhow!(\"Failed to get streaming stats from worker\"))\n                .map_err(err)?;\n\n            // Aggregate fragment_stats\n            for (fragment_id, fragment_stats) in result.fragment_stats {\n                if let Some(s) = all.fragment_stats.get_mut(&fragment_id) {\n                    s.actor_count += fragment_stats.actor_count;\n                    s.current_epoch = min(s.current_epoch, fragment_stats.current_epoch);\n                } else {\n                    all.fragment_stats.insert(fragment_id, fragment_stats);\n                }\n            }\n\n            // Aggregate relation_stats\n            for (relation_id, relation_stats) in result.relation_stats {\n                if let Some(s) = all.relation_stats.get_mut(&relation_id) {\n                    s.actor_count += relation_stats.actor_count;\n                    s.current_epoch = min(s.current_epoch, relation_stats.current_epoch);\n                } else {","sourceCodeStart":788,"sourceCodeEnd":824,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/dashboard/mod.rs#L788-L824","documentation":"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(|_| ...)`.","triggerScenarios":"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.","commonSituations":"Compute node OOM/killed during monitoring; network issues between meta and worker pods in Kubernetes; polling the dashboard mid-restart.","solutions":["Verify all workers are alive and connected (dashboard worker list / SHOW NODES) and restart unhealthy ones.","Check meta node logs for the underlying per-worker error that this message hides.","Retry after the cluster stabilizes.","Patch the handler to include the original error in the message instead of `|_|`."],"exampleFix":"// before\n.map_err(|_| anyhow!(\"Failed to get streaming stats from worker\"))\n// after\n.map_err(|e| anyhow!(\"Failed to get streaming stats from worker: {e}\"))","handlingStrategy":"retry","validationCode":"const health = await fetch('http://localhost:5691/api/v1/workers').then(r => r.json());\nif (health.some(w => w.state !== 'RUNNING')) throw new Error('A worker is not RUNNING; streaming stats from prometheus path will fail');","typeGuard":null,"tryCatchPattern":"try { const stats = await getStreamingStatsFromPrometheus(); } catch (e) {\n  if (String(e).includes('Failed to get streaming stats from worker')) {\n    // locate failing worker via meta logs, restart if needed, retry with backoff\n  }\n}","preventionTips":["Keep alerts on worker health so stats polling never runs against a degraded cluster.","Inspect meta node logs for the underlying error; the dashboard message hides it.","Throttle dashboard polling frequency in large clusters to reduce per-worker RPC pressure."],"tags":["grpc","prometheus","monitoring","risingwave"],"backgroundTag":"upstream-api-error","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"}