1use std::collections::BTreeMap;
16
17use common_meta::key::flow::flow_state::FlowStat;
18
19use crate::StreamingEngine;
20use crate::engine::FlowStatProvider;
21
22impl FlowStatProvider for StreamingEngine {
23 async fn flow_stat(&self) -> FlowStat {
24 let mut state_size_map = BTreeMap::new();
25 let mut last_exec_time_map = BTreeMap::new();
26 let mut start_time_map = BTreeMap::new();
27
28 for worker in self.worker_handles.iter() {
29 match worker.get_full_flow_stat().await {
30 Ok((sizes, exec_times, start_times)) => {
31 state_size_map.extend(sizes.into_iter().map(|(k, v)| (k as u32, v)));
32 last_exec_time_map.extend(exec_times.into_iter().map(|(k, v)| (k as u32, v)));
33 start_time_map.extend(start_times.into_iter().map(|(k, v)| (k as u32, v)));
34 }
35 Err(err) => {
36 common_telemetry::error!(err; "Get full flow stat error");
37 }
38 }
39 }
40
41 FlowStat {
42 state_size: state_size_map,
43 last_exec_time_map,
44 start_time_map,
45 }
46 }
47}