Skip to main content

polars_stream/
metrics.rs

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    /// # Panics
180    /// When debug_assertions enabled, panics if called more than once for a node within a single
181    /// phase.
182    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                // Each node should only have 1 set of metrics, identified by the Arc address.
190                // But the registration can be called multiple times (per phase).
191                assert!(Arc::ptr_eq(&io_metrics, e.get()));
192            },
193            Entry::Vacant(e) => {
194                e.insert(io_metrics);
195            },
196        };
197    }
198}