Skip to main content

promql/extension_plan/
empty_metric.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::collections::HashMap;
16use std::ops::Div;
17use std::pin::Pin;
18use std::sync::Arc;
19use std::task::{Context, Poll};
20
21use datafusion::arrow::array::ArrayRef;
22use datafusion::arrow::datatypes::{DataType, TimeUnit};
23use datafusion::catalog::Session;
24use datafusion::common::arrow::datatypes::Field;
25use datafusion::common::stats::Precision;
26use datafusion::common::tree_node::TreeNodeRecursion;
27use datafusion::common::{
28    DFSchema, DFSchemaRef, Result as DataFusionResult, Statistics, TableReference,
29};
30use datafusion::datasource::{MemTable, provider_as_source};
31use datafusion::error::DataFusionError;
32use datafusion::execution::context::TaskContext;
33use datafusion::logical_expr::physical_planning_context::PhysicalPlanningContext;
34use datafusion::logical_expr::{ExprSchemable, LogicalPlan, UserDefinedLogicalNodeCore};
35use datafusion::physical_expr::{EquivalenceProperties, PhysicalExpr, PhysicalExprRef};
36use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
37use datafusion::physical_plan::metrics::{BaselineMetrics, ExecutionPlanMetricsSet, MetricsSet};
38use datafusion::physical_plan::{
39    DisplayAs, DisplayFormatType, ExecutionPlan, Partitioning, PlanProperties, RecordBatchStream,
40    SendableRecordBatchStream, StatisticsArgs,
41};
42use datafusion::physical_planner::PhysicalPlanner;
43use datafusion::prelude::{Expr, ident, lit};
44use datafusion_expr::LogicalPlanBuilder;
45use datatypes::arrow::array::TimestampMillisecondArray;
46use datatypes::arrow::datatypes::SchemaRef;
47use datatypes::arrow::record_batch::RecordBatch;
48use futures::Stream;
49
50use crate::extension_plan::Millisecond;
51
52/// Empty source plan that generate record batch with two columns:
53/// - time index column, computed from start, end and interval
54/// - value column, generated by the input expr. The expr should not
55///   reference any column except the time index column.
56#[derive(Debug, Clone, PartialEq, Eq, Hash)]
57pub struct EmptyMetric {
58    start: Millisecond,
59    end: Millisecond,
60    interval: Millisecond,
61    expr: Option<Expr>,
62    /// Schema that only contains the time index column.
63    /// This is for intermediate result only.
64    time_index_schema: DFSchemaRef,
65    /// Schema of the output record batch
66    result_schema: DFSchemaRef,
67    // This dummy input's sole purpose is to provide a schema for use in DataFusion's
68    // `SimplifyExpressions`. Otherwise it may report a "no field name ..." error.
69    // The error is caused by an optimization that tries to rewrite "A = A", during
70    // which will find the field in plan's schema. However, the schema is empty if the
71    // plan does not have an input.
72    dummy_input: LogicalPlan,
73}
74
75impl EmptyMetric {
76    pub fn new(
77        start: Millisecond,
78        end: Millisecond,
79        interval: Millisecond,
80        time_index_column_name: String,
81        field_column_name: String,
82        field_expr: Option<Expr>,
83    ) -> DataFusionResult<Self> {
84        let qualifier = Some(TableReference::bare(""));
85        let ts_only_schema = build_ts_only_schema(&time_index_column_name);
86        let mut fields = vec![(qualifier.clone(), ts_only_schema.field(0).clone())];
87        if let Some(field_expr) = &field_expr {
88            let field_data_type = field_expr.get_type(&ts_only_schema)?;
89            fields.push((
90                qualifier.clone(),
91                Arc::new(Field::new(field_column_name, field_data_type, true)),
92            ));
93        }
94        let schema = Arc::new(DFSchema::new_with_metadata(fields, HashMap::new())?);
95
96        let table = MemTable::try_new(Arc::new(schema.as_arrow().clone()), vec![vec![]])?;
97        let source = provider_as_source(Arc::new(table));
98        let dummy_input =
99            LogicalPlanBuilder::scan("dummy", source, None).and_then(|x| x.build())?;
100
101        Ok(Self {
102            start,
103            end,
104            interval,
105            time_index_schema: Arc::new(ts_only_schema),
106            result_schema: schema,
107            expr: field_expr,
108            dummy_input,
109        })
110    }
111
112    pub const fn name() -> &'static str {
113        "EmptyMetric"
114    }
115
116    pub fn to_execution_plan(
117        &self,
118        session: &dyn Session,
119        physical_planner: &dyn PhysicalPlanner,
120        planning_ctx: &PhysicalPlanningContext,
121    ) -> DataFusionResult<Arc<dyn ExecutionPlan>> {
122        let physical_expr = self
123            .expr
124            .as_ref()
125            .map(|expr| {
126                physical_planner.create_physical_expr(
127                    expr,
128                    &self.time_index_schema,
129                    session,
130                    planning_ctx,
131                )
132            })
133            .transpose()?;
134        let result_schema: SchemaRef = self.result_schema.inner().clone();
135        let properties = Arc::new(PlanProperties::new(
136            EquivalenceProperties::new(result_schema.clone()),
137            Partitioning::UnknownPartitioning(1),
138            EmissionType::Incremental,
139            Boundedness::Bounded,
140        ));
141        Ok(Arc::new(EmptyMetricExec {
142            start: self.start,
143            end: self.end,
144            interval: self.interval,
145            time_index_schema: self.time_index_schema.inner().clone(),
146            result_schema,
147            expr: physical_expr,
148            properties,
149            metric: ExecutionPlanMetricsSet::new(),
150        }))
151    }
152}
153
154impl UserDefinedLogicalNodeCore for EmptyMetric {
155    fn name(&self) -> &str {
156        Self::name()
157    }
158
159    fn inputs(&self) -> Vec<&LogicalPlan> {
160        vec![&self.dummy_input]
161    }
162
163    fn schema(&self) -> &DFSchemaRef {
164        &self.result_schema
165    }
166
167    fn expressions(&self) -> Vec<Expr> {
168        if let Some(expr) = &self.expr {
169            vec![expr.clone()]
170        } else {
171            vec![]
172        }
173    }
174
175    fn fmt_for_explain(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
176        write!(
177            f,
178            "EmptyMetric: range=[{}..{}], interval=[{}]",
179            self.start, self.end, self.interval,
180        )
181    }
182
183    fn with_exprs_and_inputs(
184        &self,
185        exprs: Vec<Expr>,
186        _inputs: Vec<LogicalPlan>,
187    ) -> DataFusionResult<Self> {
188        Ok(Self {
189            start: self.start,
190            end: self.end,
191            interval: self.interval,
192            expr: exprs.into_iter().next(),
193            time_index_schema: self.time_index_schema.clone(),
194            result_schema: self.result_schema.clone(),
195            dummy_input: self.dummy_input.clone(),
196        })
197    }
198}
199
200impl PartialOrd for EmptyMetric {
201    fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
202        // Compare fields in order excluding schema fields
203        match self.start.partial_cmp(&other.start) {
204            Some(core::cmp::Ordering::Equal) => {}
205            ord => return ord,
206        }
207        match self.end.partial_cmp(&other.end) {
208            Some(core::cmp::Ordering::Equal) => {}
209            ord => return ord,
210        }
211        match self.interval.partial_cmp(&other.interval) {
212            Some(core::cmp::Ordering::Equal) => {}
213            ord => return ord,
214        }
215        self.expr.partial_cmp(&other.expr)
216    }
217}
218
219#[derive(Debug, Clone)]
220pub struct EmptyMetricExec {
221    start: Millisecond,
222    end: Millisecond,
223    interval: Millisecond,
224    /// Schema that only contains the time index column.
225    /// This is for intermediate result only.
226    time_index_schema: SchemaRef,
227    /// Schema of the output record batch
228    result_schema: SchemaRef,
229    expr: Option<PhysicalExprRef>,
230    properties: Arc<PlanProperties>,
231    metric: ExecutionPlanMetricsSet,
232}
233
234impl ExecutionPlan for EmptyMetricExec {
235    fn apply_expressions(
236        &self,
237        f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> datafusion_common::Result<TreeNodeRecursion>,
238    ) -> DataFusionResult<TreeNodeRecursion> {
239        datafusion::physical_plan::apply_expression_roots(self.expr.iter(), f)
240    }
241
242    fn schema(&self) -> SchemaRef {
243        self.result_schema.clone()
244    }
245
246    fn properties(&self) -> &Arc<PlanProperties> {
247        &self.properties
248    }
249
250    fn maintains_input_order(&self) -> Vec<bool> {
251        vec![]
252    }
253
254    fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
255        vec![]
256    }
257
258    fn with_new_children(
259        self: Arc<Self>,
260        _children: Vec<Arc<dyn ExecutionPlan>>,
261    ) -> DataFusionResult<Arc<dyn ExecutionPlan>> {
262        Ok(Arc::new(self.as_ref().clone()))
263    }
264
265    fn execute(
266        &self,
267        partition: usize,
268        _context: Arc<TaskContext>,
269    ) -> DataFusionResult<SendableRecordBatchStream> {
270        let baseline_metric = BaselineMetrics::new(&self.metric, partition);
271        Ok(Box::pin(EmptyMetricStream {
272            start: self.start,
273            end: self.end,
274            interval: self.interval,
275            expr: self.expr.clone(),
276            is_first_poll: true,
277            time_index_schema: self.time_index_schema.clone(),
278            result_schema: self.result_schema.clone(),
279            metric: baseline_metric,
280        }))
281    }
282
283    fn metrics(&self) -> Option<MetricsSet> {
284        Some(self.metric.clone_inner())
285    }
286
287    fn statistics_from_inputs(
288        &self,
289        _input_stats: &[Arc<Statistics>],
290        args: &StatisticsArgs,
291    ) -> DataFusionResult<Arc<Statistics>> {
292        let partition = args.partition();
293        if partition.is_some() {
294            return Ok(Arc::new(Statistics::new_unknown(self.schema().as_ref())));
295        }
296
297        let estimated_row_num = if self.end > self.start {
298            (self.end - self.start) as f64 / self.interval as f64
299        } else {
300            0.0
301        };
302        let total_byte_size = estimated_row_num * std::mem::size_of::<Millisecond>() as f64;
303
304        Ok(Arc::new(Statistics {
305            num_rows: Precision::Inexact(estimated_row_num.floor() as _),
306            total_byte_size: Precision::Inexact(total_byte_size.floor() as _),
307            column_statistics: Statistics::unknown_column(&self.schema()),
308        }))
309    }
310
311    fn name(&self) -> &str {
312        "EmptyMetricExec"
313    }
314}
315
316impl DisplayAs for EmptyMetricExec {
317    fn fmt_as(&self, t: DisplayFormatType, f: &mut std::fmt::Formatter) -> std::fmt::Result {
318        match t {
319            DisplayFormatType::Default
320            | DisplayFormatType::Verbose
321            | DisplayFormatType::TreeRender => write!(
322                f,
323                "EmptyMetric: range=[{}..{}], interval=[{}]",
324                self.start, self.end, self.interval,
325            ),
326        }
327    }
328}
329
330pub struct EmptyMetricStream {
331    start: Millisecond,
332    end: Millisecond,
333    interval: Millisecond,
334    expr: Option<PhysicalExprRef>,
335    /// This stream only generate one record batch at the first poll
336    is_first_poll: bool,
337    /// Schema that only contains the time index column.
338    /// This is for intermediate result only.
339    time_index_schema: SchemaRef,
340    /// Schema of the output record batch
341    result_schema: SchemaRef,
342    metric: BaselineMetrics,
343}
344
345impl RecordBatchStream for EmptyMetricStream {
346    fn schema(&self) -> SchemaRef {
347        self.result_schema.clone()
348    }
349}
350
351impl Stream for EmptyMetricStream {
352    type Item = DataFusionResult<RecordBatch>;
353
354    fn poll_next(mut self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
355        let result = if self.is_first_poll {
356            self.is_first_poll = false;
357            let _timer = self.metric.elapsed_compute().timer();
358
359            // build the time index array, and a record batch that
360            // only contains that array as the input of field expr
361            let time_array = (self.start..=self.end)
362                .step_by(self.interval as _)
363                .collect::<Vec<_>>();
364            let time_array = Arc::new(TimestampMillisecondArray::from(time_array));
365            let num_rows = time_array.len();
366            let input_record_batch =
367                RecordBatch::try_new(self.time_index_schema.clone(), vec![time_array.clone()])
368                    .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))?;
369            let mut result_arrays: Vec<ArrayRef> = vec![time_array];
370
371            // evaluate the field expr and get the result
372            if let Some(field_expr) = &self.expr {
373                result_arrays.push(
374                    field_expr
375                        .evaluate(&input_record_batch)
376                        .and_then(|x| x.into_array(num_rows))?,
377                );
378            }
379
380            // assemble the output record batch
381            let batch = RecordBatch::try_new(self.result_schema.clone(), result_arrays)
382                .map_err(|e| DataFusionError::ArrowError(Box::new(e), None));
383
384            Poll::Ready(Some(batch))
385        } else {
386            Poll::Ready(None)
387        };
388        self.metric.record_poll(result)
389    }
390}
391
392/// Build a schema that only contains **millisecond** timestamp column
393fn build_ts_only_schema(column_name: &str) -> DFSchema {
394    let ts_field = Field::new(
395        column_name,
396        DataType::Timestamp(TimeUnit::Millisecond, None),
397        false,
398    );
399    // safety: should not fail (UT covers this)
400    DFSchema::new_with_metadata(
401        vec![(Some(TableReference::bare("")), Arc::new(ts_field))],
402        HashMap::new(),
403    )
404    .unwrap()
405}
406
407// Convert timestamp column to UNIX epoch second:
408// https://prometheus.io/docs/prometheus/latest/querying/functions/#time
409pub fn build_special_time_expr(time_index_column_name: &str) -> Expr {
410    let input_schema = build_ts_only_schema(time_index_column_name);
411    // safety: should not failed (UT covers this)
412    ident(time_index_column_name)
413        .cast_to(&DataType::Int64, &input_schema)
414        .unwrap()
415        .cast_to(&DataType::Float64, &input_schema)
416        .unwrap()
417        .div(lit(1000.0)) // cast to second will lost precision, so we cast to float64 first and manually divide by 1000
418}
419
420#[cfg(test)]
421mod test {
422    use datafusion::physical_planner::DefaultPhysicalPlanner;
423    use datafusion::prelude::SessionContext;
424
425    use super::*;
426
427    async fn do_empty_metric_test(
428        start: Millisecond,
429        end: Millisecond,
430        interval: Millisecond,
431        time_column_name: String,
432        field_column_name: String,
433        expected: String,
434    ) {
435        let session_context = SessionContext::default();
436        let df_default_physical_planner = DefaultPhysicalPlanner::default();
437        let time_expr = build_special_time_expr(&time_column_name);
438        let empty_metric = EmptyMetric::new(
439            start,
440            end,
441            interval,
442            time_column_name,
443            field_column_name,
444            Some(time_expr),
445        )
446        .unwrap();
447        let empty_metric_exec = empty_metric
448            .to_execution_plan(
449                &session_context.state(),
450                &df_default_physical_planner,
451                &PhysicalPlanningContext::default(),
452            )
453            .unwrap();
454
455        let result =
456            datafusion::physical_plan::collect(empty_metric_exec, session_context.task_ctx())
457                .await
458                .unwrap();
459        let result_literal = datatypes::arrow::util::pretty::pretty_format_batches(&result)
460            .unwrap()
461            .to_string();
462
463        assert_eq!(result_literal, expected);
464    }
465
466    #[tokio::test]
467    async fn normal_empty_metric_test() {
468        do_empty_metric_test(
469            0,
470            100,
471            10,
472            "time".to_string(),
473            "value".to_string(),
474            String::from(
475                "+-------------------------+-------+\
476                \n| time                    | value |\
477                \n+-------------------------+-------+\
478                \n| 1970-01-01T00:00:00     | 0.0   |\
479                \n| 1970-01-01T00:00:00.010 | 0.01  |\
480                \n| 1970-01-01T00:00:00.020 | 0.02  |\
481                \n| 1970-01-01T00:00:00.030 | 0.03  |\
482                \n| 1970-01-01T00:00:00.040 | 0.04  |\
483                \n| 1970-01-01T00:00:00.050 | 0.05  |\
484                \n| 1970-01-01T00:00:00.060 | 0.06  |\
485                \n| 1970-01-01T00:00:00.070 | 0.07  |\
486                \n| 1970-01-01T00:00:00.080 | 0.08  |\
487                \n| 1970-01-01T00:00:00.090 | 0.09  |\
488                \n| 1970-01-01T00:00:00.100 | 0.1   |\
489                \n+-------------------------+-------+",
490            ),
491        )
492        .await
493    }
494
495    #[tokio::test]
496    async fn unaligned_empty_metric_test() {
497        do_empty_metric_test(
498            0,
499            100,
500            11,
501            "time".to_string(),
502            "value".to_string(),
503            String::from(
504                "+-------------------------+-------+\
505                \n| time                    | value |\
506                \n+-------------------------+-------+\
507                \n| 1970-01-01T00:00:00     | 0.0   |\
508                \n| 1970-01-01T00:00:00.011 | 0.011 |\
509                \n| 1970-01-01T00:00:00.022 | 0.022 |\
510                \n| 1970-01-01T00:00:00.033 | 0.033 |\
511                \n| 1970-01-01T00:00:00.044 | 0.044 |\
512                \n| 1970-01-01T00:00:00.055 | 0.055 |\
513                \n| 1970-01-01T00:00:00.066 | 0.066 |\
514                \n| 1970-01-01T00:00:00.077 | 0.077 |\
515                \n| 1970-01-01T00:00:00.088 | 0.088 |\
516                \n| 1970-01-01T00:00:00.099 | 0.099 |\
517                \n+-------------------------+-------+",
518            ),
519        )
520        .await
521    }
522
523    #[tokio::test]
524    async fn one_row_empty_metric_test() {
525        do_empty_metric_test(
526            0,
527            100,
528            1000,
529            "time".to_string(),
530            "value".to_string(),
531            String::from(
532                "+---------------------+-------+\
533                \n| time                | value |\
534                \n+---------------------+-------+\
535                \n| 1970-01-01T00:00:00 | 0.0   |\
536                \n+---------------------+-------+",
537            ),
538        )
539        .await
540    }
541
542    #[tokio::test]
543    async fn negative_range_empty_metric_test() {
544        do_empty_metric_test(
545            1000,
546            -1000,
547            10,
548            "time".to_string(),
549            "value".to_string(),
550            String::from(
551                "+------+-------+\
552                \n| time | value |\
553                \n+------+-------+\
554                \n+------+-------+",
555            ),
556        )
557        .await
558    }
559
560    #[tokio::test]
561    async fn no_field_expr() {
562        let session_context = SessionContext::default();
563        let df_default_physical_planner = DefaultPhysicalPlanner::default();
564        let empty_metric =
565            EmptyMetric::new(0, 200, 1000, "time".to_string(), "value".to_string(), None).unwrap();
566        let empty_metric_exec = empty_metric
567            .to_execution_plan(
568                &session_context.state(),
569                &df_default_physical_planner,
570                &PhysicalPlanningContext::default(),
571            )
572            .unwrap();
573
574        let result =
575            datafusion::physical_plan::collect(empty_metric_exec, session_context.task_ctx())
576                .await
577                .unwrap();
578        let result_literal = datatypes::arrow::util::pretty::pretty_format_batches(&result)
579            .unwrap()
580            .to_string();
581
582        let expected = String::from(
583            "+---------------------+\
584            \n| time                |\
585            \n+---------------------+\
586            \n| 1970-01-01T00:00:00 |\
587            \n+---------------------+",
588        );
589        assert_eq!(result_literal, expected);
590    }
591}