1use std::sync::Arc;
2use std::time::Duration;
3
4use polars_async::executor::TaskMetrics;
5pub use polars_io::metrics::{IOMetrics, OptIOMetrics};
6use slotmap::{SecondaryMap, SlotMap};
7
8use crate::LogicalPipe;
9use crate::graph::{GraphNodeKey, LogicalPipeKey};
10use crate::pipe::PipeMetrics;
11
12#[derive(Default, Clone)]
13pub struct NodeMetrics {
14 pub total_polls: u64,
15 pub total_stolen_polls: u64,
16 pub total_poll_time_ns: u64,
17 pub max_poll_time_ns: u64,
18
19 pub total_state_updates: u64,
20 pub total_state_update_time_ns: u64,
21 pub max_state_update_time_ns: u64,
22
23 pub morsels_sent: u64,
24 pub rows_sent: u64,
25 pub largest_morsel_sent: u64,
26 pub morsels_received: u64,
27 pub rows_received: u64,
28 pub largest_morsel_received: u64,
29
30 pub io_total_active_ns: u64,
31 pub io_total_bytes_requested: u64,
32 pub io_total_bytes_received: u64,
33 pub io_total_bytes_sent: u64,
34
35 pub state_update_in_progress: bool,
36 pub num_running_tasks: u32,
37 pub done: bool,
38}
39
40impl NodeMetrics {
41 fn add_task(&mut self, task_metrics: &TaskMetrics) {
42 self.total_polls += task_metrics.total_polls.load();
43 self.total_stolen_polls += task_metrics.total_stolen_polls.load();
44 self.total_poll_time_ns += task_metrics.total_poll_time_ns.load();
45 self.max_poll_time_ns = self
46 .max_poll_time_ns
47 .max(task_metrics.max_poll_time_ns.load());
48 self.num_running_tasks += (!task_metrics.done.load()) as u32;
49 }
50
51 fn add_io(&mut self, io_metrics: &IOMetrics) {
52 self.io_total_active_ns += io_metrics.io_timer.total_time_live_ns();
53 self.io_total_bytes_requested += io_metrics.bytes_requested.load();
54 self.io_total_bytes_received += io_metrics.bytes_received.load();
55 self.io_total_bytes_sent += io_metrics.bytes_sent.load();
56 }
57
58 fn reset_io_metrics(&mut self) {
59 self.io_total_active_ns = 0;
60 self.io_total_bytes_requested = 0;
61 self.io_total_bytes_received = 0;
62 self.io_total_bytes_sent = 0;
63 }
64
65 fn start_state_update(&mut self) {
66 self.state_update_in_progress = true;
67 }
68
69 fn stop_state_update(&mut self, time: Duration, is_done: bool) {
70 let time_ns = time.as_nanos() as u64;
71 self.total_state_updates += 1;
72 self.total_state_update_time_ns += time_ns;
73 self.max_state_update_time_ns = self.max_state_update_time_ns.max(time_ns);
74 self.state_update_in_progress = false;
75 self.done = is_done;
76 }
77
78 fn add_send_metrics(&mut self, pipe_metrics: &PipeMetrics) {
79 self.morsels_sent += pipe_metrics.morsels_sent.load();
80 self.rows_sent += pipe_metrics.rows_sent.load();
81 self.largest_morsel_sent = self
82 .largest_morsel_sent
83 .max(pipe_metrics.largest_morsel_sent.load());
84 }
85
86 fn add_recv_metrics(&mut self, pipe_metrics: &PipeMetrics) {
87 self.morsels_received += pipe_metrics.morsels_received.load();
88 self.rows_received += pipe_metrics.rows_received.load();
89 self.largest_morsel_received = self
90 .largest_morsel_received
91 .max(pipe_metrics.largest_morsel_received.load());
92 }
93}
94
95#[derive(Default, Clone)]
96pub struct GraphMetrics {
97 node_metrics: SecondaryMap<GraphNodeKey, NodeMetrics>,
98 in_progress_io_metrics: SecondaryMap<GraphNodeKey, Arc<IOMetrics>>,
99 in_progress_task_metrics: SecondaryMap<GraphNodeKey, Vec<Arc<TaskMetrics>>>,
100 in_progress_pipe_metrics: SecondaryMap<LogicalPipeKey, Vec<Arc<PipeMetrics>>>,
101}
102
103impl GraphMetrics {
104 pub fn add_task(&mut self, key: GraphNodeKey, task_metrics: Arc<TaskMetrics>) {
105 self.in_progress_task_metrics
106 .entry(key)
107 .unwrap()
108 .or_default()
109 .push(task_metrics);
110 }
111
112 pub fn add_pipe(&mut self, key: LogicalPipeKey, pipe_metrics: Arc<PipeMetrics>) {
113 self.in_progress_pipe_metrics
114 .entry(key)
115 .unwrap()
116 .or_default()
117 .push(pipe_metrics);
118 }
119
120 pub fn start_state_update(&mut self, key: GraphNodeKey) {
121 self.node_metrics
122 .entry(key)
123 .unwrap()
124 .or_default()
125 .start_state_update();
126 }
127
128 pub fn stop_state_update(&mut self, key: GraphNodeKey, time: Duration, is_done: bool) {
129 self.node_metrics[key].stop_state_update(time, is_done);
130 }
131
132 pub fn flush(&mut self, pipes: &SlotMap<LogicalPipeKey, LogicalPipe>) {
133 for (key, in_progress_task_metrics) in self.in_progress_task_metrics.iter_mut() {
134 let this_node_metrics = self.node_metrics.entry(key).unwrap().or_default();
135 this_node_metrics.num_running_tasks = 0;
136 for task_metrics in in_progress_task_metrics.drain(..) {
137 this_node_metrics.add_task(&task_metrics);
138 }
139 }
140
141 for (key, io_metrics) in self.in_progress_io_metrics.iter_mut() {
142 let this_node_metrics = self.node_metrics.entry(key).unwrap().or_default();
143 this_node_metrics.reset_io_metrics();
144 this_node_metrics.add_io(io_metrics);
145 }
146
147 for (key, in_progress_pipe_metrics) in self.in_progress_pipe_metrics.iter_mut() {
148 for pipe_metrics in in_progress_pipe_metrics.drain(..) {
149 let pipe = &pipes[key];
150 self.node_metrics
151 .entry(pipe.receiver)
152 .unwrap()
153 .or_default()
154 .add_recv_metrics(&pipe_metrics);
155 self.node_metrics
156 .entry(pipe.sender)
157 .unwrap()
158 .or_default()
159 .add_send_metrics(&pipe_metrics);
160 }
161 }
162 }
163
164 pub fn get(&self, key: GraphNodeKey) -> Option<&NodeMetrics> {
165 self.node_metrics.get(key)
166 }
167
168 pub fn iter(&self) -> slotmap::secondary::Iter<'_, GraphNodeKey, NodeMetrics> {
169 self.node_metrics.iter()
170 }
171}
172
173pub struct NodeMetricsRegistrator {
174 pub graph_key: GraphNodeKey,
175 pub graph_metrics: Arc<parking_lot::Mutex<GraphMetrics>>,
176}
177
178impl NodeMetricsRegistrator {
179 pub fn register_io_metrics(&self, io_metrics: Arc<IOMetrics>) {
183 let mut guard = self.graph_metrics.lock();
184
185 use slotmap::secondary::Entry;
186
187 match guard.in_progress_io_metrics.entry(self.graph_key).unwrap() {
188 Entry::Occupied(e) => {
189 assert!(Arc::ptr_eq(&io_metrics, e.get()));
192 },
193 Entry::Vacant(e) => {
194 e.insert(io_metrics);
195 },
196 };
197 }
198}