1use std::time::Duration;
16
17use datafusion::physical_plan::metrics::{
18 Count, ExecutionPlanMetricsSet, MetricBuilder, ScopedTimerGuard, Time, Timestamp,
19};
20
21#[derive(Debug, Clone)]
25pub struct StreamMetrics {
26 end_time: Timestamp,
28 output_rows: Count,
30 output_bytes: Count,
32 poll_elapsed: Time,
34 await_elapsed: Time,
36 dynamic_filter_cost: Time,
38}
39
40impl StreamMetrics {
41 pub fn new(metrics: &ExecutionPlanMetricsSet, partition: usize) -> Self {
43 let start_time = MetricBuilder::new(metrics).start_timestamp(partition);
44 start_time.record();
45
46 Self {
47 end_time: MetricBuilder::new(metrics).end_timestamp(partition),
48 output_rows: MetricBuilder::new(metrics).output_rows(partition),
49 output_bytes: MetricBuilder::new(metrics).output_bytes(partition),
50 poll_elapsed: MetricBuilder::new(metrics).subset_time("elapsed_poll", partition),
51 await_elapsed: MetricBuilder::new(metrics).subset_time("elapsed_await", partition),
52 dynamic_filter_cost: MetricBuilder::new(metrics)
53 .subset_time("dynamic_filter_cost", partition),
54 }
55 }
56
57 pub fn record_output(&self, num_rows: usize) {
58 self.output_rows.add(num_rows);
59 }
60
61 pub fn record_output_bytes(&self, num_bytes: usize) {
62 self.output_bytes.add(num_bytes);
63 }
64
65 pub fn try_done(&self) {
67 if self.end_time.value().is_none() {
68 self.end_time.record()
69 }
70 }
71
72 pub fn poll_timer(&self) -> ScopedTimerGuard<'_> {
74 self.poll_elapsed.timer()
75 }
76
77 pub fn dynamic_filter_timer(&self) -> ScopedTimerGuard<'_> {
78 self.dynamic_filter_cost.timer()
79 }
80
81 pub fn record_await_duration(&self, duration: Duration) {
82 self.await_elapsed.add_duration(duration);
83 }
84}
85
86impl Drop for StreamMetrics {
87 fn drop(&mut self) {
88 self.try_done()
89 }
90}