Skip to main content

table/table/
metrics.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::time::Duration;
16
17use datafusion::physical_plan::metrics::{
18    Count, ExecutionPlanMetricsSet, MetricBuilder, ScopedTimerGuard, Time, Timestamp,
19};
20
21/// This metrics struct is used to record and hold metrics like memory usage
22/// of result batch in [`crate::table::scan::StreamWithMetricWrapper`]
23/// during query execution.
24#[derive(Debug, Clone)]
25pub struct StreamMetrics {
26    /// Timestamp when the stream finished
27    end_time: Timestamp,
28    /// Number of rows in output
29    output_rows: Count,
30    /// Number of bytes in output
31    output_bytes: Count,
32    /// Elapsed time used to `poll` the stream
33    poll_elapsed: Time,
34    /// Elapsed time used to `.await`ing the stream
35    await_elapsed: Time,
36    /// Time spent evaluating and applying decoded dynamic filters.
37    dynamic_filter_cost: Time,
38}
39
40impl StreamMetrics {
41    /// Create a new [`StreamMetrics`] structure, and set `start_time` to now.
42    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    /// Record the end time of the query
66    pub fn try_done(&self) {
67        if self.end_time.value().is_none() {
68            self.end_time.record()
69        }
70    }
71
72    /// Return a timer guard that records the time elapsed in poll
73    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}