Skip to main content

query/range_select/
plan.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::cmp::Ordering;
16use std::collections::btree_map::Entry;
17use std::collections::{BTreeMap, HashMap};
18use std::fmt::Display;
19use std::pin::Pin;
20use std::sync::Arc;
21use std::task::{Context, Poll};
22use std::time::Duration;
23
24use arrow::compute::{self, CastOptions, cast_with_options, take_arrays};
25use arrow_schema::{DataType, Field, Schema, SchemaRef, SortOptions, TimeUnit};
26use common_function::aggrs::aggr_wrapper::get_aggr_func;
27use common_recordbatch::DfSendableRecordBatchStream;
28use datafusion::catalog::Session;
29use datafusion::common::Result as DataFusionResult;
30use datafusion::error::Result as DfResult;
31use datafusion::execution::TaskContext;
32use datafusion::logical_expr::physical_planning_context::PhysicalPlanningContext;
33use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
34use datafusion::physical_plan::metrics::{BaselineMetrics, ExecutionPlanMetricsSet, MetricsSet};
35use datafusion::physical_plan::{
36    DisplayAs, DisplayFormatType, ExecutionPlan, PlanProperties, RecordBatchStream,
37    SendableRecordBatchStream, apply_expression_roots,
38};
39use datafusion_common::hash_utils::{RandomState, create_hashes};
40use datafusion_common::tree_node::TreeNodeRecursion;
41use datafusion_common::{DFSchema, DFSchemaRef, DataFusionError, ScalarValue};
42use datafusion_expr::utils::{COUNT_STAR_EXPANSION, exprlist_to_fields};
43use datafusion_expr::{
44    Accumulator, Expr, ExprSchemable, LogicalPlan, UserDefinedLogicalNodeCore, lit,
45};
46use datafusion_physical_expr::aggregate::{AggregateExprBuilder, AggregateFunctionExpr};
47use datafusion_physical_expr::{
48    Distribution, EquivalenceProperties, Partitioning, PhysicalExpr, PhysicalSortExpr,
49    create_physical_expr, create_physical_sort_expr,
50};
51use datatypes::arrow::array::{
52    Array, ArrayRef, TimestampMillisecondArray, TimestampMillisecondBuilder, UInt32Builder,
53};
54use datatypes::arrow::datatypes::{ArrowPrimitiveType, TimestampMillisecondType};
55use datatypes::arrow::record_batch::RecordBatch;
56use datatypes::arrow::row::{OwnedRow, RowConverter, SortField};
57use futures::{Stream, ready};
58use futures_util::StreamExt;
59use snafu::ensure;
60
61use crate::error::{RangeQuerySnafu, Result};
62
63type Millisecond = <TimestampMillisecondType as ArrowPrimitiveType>::Native;
64
65#[derive(PartialEq, Eq, Debug, Hash, Clone)]
66pub enum Fill {
67    Null,
68    Prev,
69    Linear,
70    Const(ScalarValue),
71}
72
73impl Display for Fill {
74    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
75        match self {
76            Fill::Null => write!(f, "NULL"),
77            Fill::Prev => write!(f, "PREV"),
78            Fill::Linear => write!(f, "LINEAR"),
79            Fill::Const(x) => write!(f, "{}", x),
80        }
81    }
82}
83
84impl Fill {
85    pub fn try_from_str(value: &str, datatype: &DataType) -> DfResult<Option<Self>> {
86        let s = value.to_uppercase();
87        match s.as_str() {
88            "" => Ok(None),
89            "NULL" => Ok(Some(Self::Null)),
90            "PREV" => Ok(Some(Self::Prev)),
91            "LINEAR" => {
92                if datatype.is_numeric() {
93                    Ok(Some(Self::Linear))
94                } else {
95                    Err(DataFusionError::Plan(format!(
96                        "Use FILL LINEAR on Non-numeric DataType {}",
97                        datatype
98                    )))
99                }
100            }
101            _ => ScalarValue::try_from_string(s.clone(), datatype)
102                .map_err(|err| {
103                    DataFusionError::Plan(format!(
104                        "{} is not a valid fill option, fail to convert to a const value. {{ {} }}",
105                        s, err
106                    ))
107                })
108                .map(|x| Some(Fill::Const(x))),
109        }
110    }
111
112    /// The input `data` contains data on a complete time series.
113    /// If the filling strategy is `PREV` or `LINEAR`, caller must be ensured that the incoming `ts`&`data` is ascending time order.
114    pub fn apply_fill_strategy(&self, ts: &[i64], data: &mut [ScalarValue]) -> DfResult<()> {
115        // No calculation need in `Fill::Null`
116        if matches!(self, Fill::Null) {
117            return Ok(());
118        }
119        let len = data.len();
120        if *self == Fill::Linear {
121            return Self::fill_linear(ts, data);
122        }
123        for i in 0..len {
124            if data[i].is_null() {
125                match self {
126                    Fill::Prev => {
127                        if i != 0 {
128                            data[i] = data[i - 1].clone()
129                        }
130                    }
131                    // The calculation of linear interpolation is relatively complicated.
132                    // `Self::fill_linear` is used to dispose `Fill::Linear`.
133                    // No calculation need in `Fill::Null`
134                    Fill::Linear | Fill::Null => unreachable!(),
135                    Fill::Const(v) => data[i] = v.clone(),
136                }
137            }
138        }
139        Ok(())
140    }
141
142    fn fill_linear(ts: &[i64], data: &mut [ScalarValue]) -> DfResult<()> {
143        let not_null_num = data
144            .iter()
145            .fold(0, |acc, x| if x.is_null() { acc } else { acc + 1 });
146        // We need at least two non-empty data points to perform linear interpolation
147        if not_null_num < 2 {
148            return Ok(());
149        }
150        let mut index = 0;
151        let mut head: Option<usize> = None;
152        let mut tail: Option<usize> = None;
153        while index < data.len() {
154            // find null interval [start, end)
155            // start is null, end is not-null
156            let start = data[index..]
157                .iter()
158                .position(ScalarValue::is_null)
159                .unwrap_or(data.len() - index)
160                + index;
161            if start == data.len() {
162                break;
163            }
164            let end = data[start..]
165                .iter()
166                .position(|r| !r.is_null())
167                .unwrap_or(data.len() - start)
168                + start;
169            index = end + 1;
170            // head or tail null dispose later, record start/end first
171            if start == 0 {
172                head = Some(end);
173            } else if end == data.len() {
174                tail = Some(start);
175            } else {
176                linear_interpolation(ts, data, start - 1, end, start, end)?;
177            }
178        }
179        // dispose head null interval
180        if let Some(end) = head {
181            linear_interpolation(ts, data, end, end + 1, 0, end)?;
182        }
183        // dispose tail null interval
184        if let Some(start) = tail {
185            linear_interpolation(ts, data, start - 2, start - 1, start, data.len())?;
186        }
187        Ok(())
188    }
189}
190
191/// use `(ts[i1], data[i1])`, `(ts[i2], data[i2])` as endpoint, linearly interpolates element over the interval `[start, end)`
192fn linear_interpolation(
193    ts: &[i64],
194    data: &mut [ScalarValue],
195    i1: usize,
196    i2: usize,
197    start: usize,
198    end: usize,
199) -> DfResult<()> {
200    let (x0, x1) = (ts[i1] as f64, ts[i2] as f64);
201    let (y0, y1, is_float32) = match (&data[i1], &data[i2]) {
202        (ScalarValue::Float64(Some(y0)), ScalarValue::Float64(Some(y1))) => (*y0, *y1, false),
203        (ScalarValue::Float32(Some(y0)), ScalarValue::Float32(Some(y1))) => {
204            (*y0 as f64, *y1 as f64, true)
205        }
206        _ => {
207            return Err(DataFusionError::Execution(
208                "RangePlan: Apply Fill LINEAR strategy on Non-floating type".to_string(),
209            ));
210        }
211    };
212    // To avoid divide zero error, kind of defensive programming
213    if x1 == x0 {
214        return Err(DataFusionError::Execution(
215            "RangePlan: Linear interpolation using the same coordinate points".to_string(),
216        ));
217    }
218    for i in start..end {
219        let val = y0 + (y1 - y0) / (x1 - x0) * (ts[i] as f64 - x0);
220        data[i] = if is_float32 {
221            ScalarValue::Float32(Some(val as f32))
222        } else {
223            ScalarValue::Float64(Some(val))
224        }
225    }
226    Ok(())
227}
228
229#[derive(Eq, Clone, Debug)]
230pub struct RangeFn {
231    /// with format like `max(a) RANGE 300s [FILL NULL]`
232    pub name: String,
233    pub data_type: DataType,
234    pub expr: Expr,
235    pub range: Duration,
236    pub fill: Option<Fill>,
237    /// If the `FIll` strategy is `Linear` and the output is an integer,
238    /// it is possible to calculate a floating point number.
239    /// So for `FILL==LINEAR`, the entire data will be implicitly converted to Float type
240    /// If `need_cast==true`, `data_type` may not consist with type `expr` generated.
241    pub need_cast: bool,
242}
243
244impl PartialEq for RangeFn {
245    fn eq(&self, other: &Self) -> bool {
246        self.name == other.name
247    }
248}
249
250impl PartialOrd for RangeFn {
251    fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
252        Some(self.cmp(other))
253    }
254}
255
256impl Ord for RangeFn {
257    fn cmp(&self, other: &Self) -> std::cmp::Ordering {
258        self.name.cmp(&other.name)
259    }
260}
261
262impl std::hash::Hash for RangeFn {
263    fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
264        self.name.hash(state);
265    }
266}
267
268impl Display for RangeFn {
269    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
270        write!(f, "{}", self.name)
271    }
272}
273
274#[derive(Debug, PartialEq, Eq, Hash)]
275pub struct RangeSelect {
276    /// The incoming logical plan
277    pub input: Arc<LogicalPlan>,
278    /// all range expressions
279    pub range_expr: Vec<RangeFn>,
280    pub align: Duration,
281    pub align_to: i64,
282    pub time_index: String,
283    pub time_expr: Expr,
284    pub by: Vec<Expr>,
285    pub schema: DFSchemaRef,
286    pub by_schema: DFSchemaRef,
287    /// If the `schema` of the `RangeSelect` happens to be the same as the content of the upper-level Projection Plan,
288    /// the final output needs to be `project` through `schema_project`,
289    /// so that we can omit the upper-level Projection Plan.
290    pub schema_project: Option<Vec<usize>>,
291    /// The schema before run projection, follow the order of `range expr | time index | by columns`
292    /// `schema_before_project  ----  schema_project ----> schema`
293    /// if `schema_project==None` then `schema_before_project==schema`
294    pub schema_before_project: DFSchemaRef,
295}
296
297impl PartialOrd for RangeSelect {
298    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
299        // Compare fields in order excluding `schema`, `by_schema`, `schema_before_project`.
300        match self.input.partial_cmp(&other.input) {
301            Some(Ordering::Equal) => {}
302            ord => return ord,
303        }
304        match self.range_expr.partial_cmp(&other.range_expr) {
305            Some(Ordering::Equal) => {}
306            ord => return ord,
307        }
308        match self.align.partial_cmp(&other.align) {
309            Some(Ordering::Equal) => {}
310            ord => return ord,
311        }
312        match self.align_to.partial_cmp(&other.align_to) {
313            Some(Ordering::Equal) => {}
314            ord => return ord,
315        }
316        match self.time_index.partial_cmp(&other.time_index) {
317            Some(Ordering::Equal) => {}
318            ord => return ord,
319        }
320        match self.time_expr.partial_cmp(&other.time_expr) {
321            Some(Ordering::Equal) => {}
322            ord => return ord,
323        }
324        match self.by.partial_cmp(&other.by) {
325            Some(Ordering::Equal) => {}
326            ord => return ord,
327        }
328        self.schema_project.partial_cmp(&other.schema_project)
329    }
330}
331
332impl RangeSelect {
333    pub fn try_new(
334        input: Arc<LogicalPlan>,
335        range_expr: Vec<RangeFn>,
336        align: Duration,
337        align_to: i64,
338        time_index: Expr,
339        by: Vec<Expr>,
340        projection_expr: &[Expr],
341    ) -> Result<Self> {
342        ensure!(
343            align.as_millis() != 0,
344            RangeQuerySnafu {
345                msg: "Can't use 0 as align in Range Query"
346            }
347        );
348        for expr in &range_expr {
349            ensure!(
350                expr.range.as_millis() != 0,
351                RangeQuerySnafu {
352                    msg: format!(
353                        "Invalid Range expr `{}`, Can't use 0 as range in Range Query",
354                        expr.name
355                    )
356                }
357            );
358        }
359        let mut fields = range_expr
360            .iter()
361            .map(
362                |RangeFn {
363                     name,
364                     data_type,
365                     fill,
366                     ..
367                 }| {
368                    let field = Field::new(
369                        name,
370                        data_type.clone(),
371                        // Only when data fill with Const option, the data can't be null
372                        !matches!(fill, Some(Fill::Const(..))),
373                    );
374                    Ok((None, Arc::new(field)))
375                },
376            )
377            .collect::<DfResult<Vec<_>>>()?;
378        // add align_ts
379        let ts_field = time_index.to_field(input.schema().as_ref())?;
380        let time_index_name = ts_field.1.name().clone();
381        fields.push(ts_field);
382        // add by
383        let by_fields = exprlist_to_fields(&by, &input)?;
384        fields.extend(by_fields.clone());
385        let schema_before_project = Arc::new(DFSchema::new_with_metadata(
386            fields,
387            input.schema().metadata().clone(),
388        )?);
389        let by_schema = Arc::new(DFSchema::new_with_metadata(
390            by_fields,
391            input.schema().metadata().clone(),
392        )?);
393        // If the results of project plan can be obtained directly from range plan without any additional
394        // calculations, no project plan is required. We can simply project the final output of the range
395        // plan to produce the final result.
396        let schema_project = projection_expr
397            .iter()
398            .map(|project_expr| {
399                if let Expr::Column(column) = project_expr {
400                    schema_before_project
401                        .index_of_column_by_name(column.relation.as_ref(), &column.name)
402                        .ok_or(())
403                } else {
404                    let (qualifier, field) = project_expr
405                        .to_field(input.schema().as_ref())
406                        .map_err(|_| ())?;
407                    schema_before_project
408                        .index_of_column_by_name(qualifier.as_ref(), field.name())
409                        .ok_or(())
410                }
411            })
412            .collect::<std::result::Result<Vec<usize>, ()>>()
413            .ok();
414        let schema = if let Some(project) = &schema_project {
415            let project_field = project
416                .iter()
417                .map(|i| {
418                    let f = schema_before_project.qualified_field(*i);
419                    (f.0.cloned(), f.1.clone())
420                })
421                .collect();
422            Arc::new(DFSchema::new_with_metadata(
423                project_field,
424                input.schema().metadata().clone(),
425            )?)
426        } else {
427            schema_before_project.clone()
428        };
429        Ok(Self {
430            input,
431            range_expr,
432            align,
433            align_to,
434            time_index: time_index_name,
435            time_expr: time_index,
436            schema,
437            by_schema,
438            by,
439            schema_project,
440            schema_before_project,
441        })
442    }
443}
444
445impl UserDefinedLogicalNodeCore for RangeSelect {
446    fn name(&self) -> &str {
447        "RangeSelect"
448    }
449
450    fn inputs(&self) -> Vec<&LogicalPlan> {
451        vec![&self.input]
452    }
453
454    fn schema(&self) -> &DFSchemaRef {
455        &self.schema
456    }
457
458    fn expressions(&self) -> Vec<Expr> {
459        self.range_expr
460            .iter()
461            .map(|expr| expr.expr.clone())
462            .chain([self.time_expr.clone()])
463            .chain(self.by.clone())
464            .collect()
465    }
466
467    fn fmt_for_explain(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
468        write!(
469            f,
470            "RangeSelect: range_exprs=[{}], align={}ms, align_to={}ms, align_by=[{}], time_index={}",
471            self.range_expr
472                .iter()
473                .map(ToString::to_string)
474                .collect::<Vec<_>>()
475                .join(", "),
476            self.align.as_millis(),
477            self.align_to,
478            self.by
479                .iter()
480                .map(ToString::to_string)
481                .collect::<Vec<_>>()
482                .join(", "),
483            self.time_index
484        )
485    }
486
487    fn with_exprs_and_inputs(
488        &self,
489        exprs: Vec<Expr>,
490        inputs: Vec<LogicalPlan>,
491    ) -> DataFusionResult<Self> {
492        if inputs.is_empty() {
493            return Err(DataFusionError::Plan(
494                "RangeSelect: inputs is empty".to_string(),
495            ));
496        }
497        if exprs.len() != self.range_expr.len() + self.by.len() + 1 {
498            return Err(DataFusionError::Plan(
499                "RangeSelect: exprs length not match".to_string(),
500            ));
501        }
502
503        let range_expr = exprs
504            .iter()
505            .zip(self.range_expr.iter())
506            .map(|(e, range)| RangeFn {
507                name: range.name.clone(),
508                data_type: range.data_type.clone(),
509                expr: e.clone(),
510                range: range.range,
511                fill: range.fill.clone(),
512                need_cast: range.need_cast,
513            })
514            .collect();
515        let time_expr = exprs[self.range_expr.len()].clone();
516        let by = exprs[self.range_expr.len() + 1..].to_vec();
517        Ok(Self {
518            align: self.align,
519            align_to: self.align_to,
520            range_expr,
521            input: Arc::new(inputs[0].clone()),
522            time_index: self.time_index.clone(),
523            time_expr,
524            schema: self.schema.clone(),
525            by,
526            by_schema: self.by_schema.clone(),
527            schema_project: self.schema_project.clone(),
528            schema_before_project: self.schema_before_project.clone(),
529        })
530    }
531}
532
533impl RangeSelect {
534    fn create_physical_expr_list(
535        &self,
536        is_count_aggr: bool,
537        exprs: &[Expr],
538        df_schema: &Arc<DFSchema>,
539        session: &dyn Session,
540        planning_ctx: &PhysicalPlanningContext,
541    ) -> DfResult<Vec<Arc<dyn PhysicalExpr>>> {
542        exprs
543            .iter()
544            .map(|e| match e {
545                // `count(*)` will be rewritten by `CountWildcardRule` into `count(1)` when optimizing logical plan.
546                // The modification occurs after range plan rewrite.
547                // At this time, aggregate plan has been replaced by a custom range plan,
548                // so `CountWildcardRule` has not been applied.
549                // We manually modify it when creating the physical plan.
550                #[expect(deprecated)]
551                Expr::Wildcard { .. } if is_count_aggr => create_physical_expr(
552                    &lit(COUNT_STAR_EXPANSION),
553                    df_schema.as_ref(),
554                    session.execution_props(),
555                    planning_ctx,
556                ),
557                _ => create_physical_expr(
558                    e,
559                    df_schema.as_ref(),
560                    session.execution_props(),
561                    planning_ctx,
562                ),
563            })
564            .collect::<DfResult<Vec<_>>>()
565    }
566
567    pub fn to_execution_plan(
568        &self,
569        logical_input: &LogicalPlan,
570        exec_input: Arc<dyn ExecutionPlan>,
571        session: &dyn Session,
572        planning_ctx: &PhysicalPlanningContext,
573    ) -> DfResult<Arc<dyn ExecutionPlan>> {
574        let fields: Vec<_> = self
575            .schema_before_project
576            .fields()
577            .iter()
578            .map(|field| Field::new(field.name(), field.data_type().clone(), field.is_nullable()))
579            .collect();
580        let by_fields: Vec<_> = self
581            .by_schema
582            .fields()
583            .iter()
584            .map(|field| Field::new(field.name(), field.data_type().clone(), field.is_nullable()))
585            .collect();
586        let input_dfschema = logical_input.schema();
587        let input_schema = exec_input.schema();
588        let range_exec: Vec<RangeFnExec> = self
589            .range_expr
590            .iter()
591            .map(|range_fn| {
592                let name = range_fn.expr.schema_name().to_string();
593                let range_expr = match &range_fn.expr {
594                    Expr::Alias(expr) => expr.expr.as_ref(),
595                    others => others,
596                };
597
598                let expr = match get_aggr_func(range_expr) {
599                    Some(aggr)
600                        if (aggr.func.name() == "last_value"
601                            || aggr.func.name() == "first_value") =>
602                    {
603                        let order_by = if !aggr.params.order_by.is_empty() {
604                            aggr.params
605                                .order_by
606                                .iter()
607                                .map(|x| {
608                                    create_physical_sort_expr(
609                                        x,
610                                        input_dfschema.as_ref(),
611                                        session.execution_props(),
612                                        planning_ctx,
613                                    )
614                                })
615                                .collect::<DfResult<Vec<_>>>()?
616                        } else {
617                            // if user not assign order by, time index is needed as default ordering
618                            let time_index = create_physical_expr(
619                                &self.time_expr,
620                                input_dfschema.as_ref(),
621                                session.execution_props(),
622                                planning_ctx,
623                            )?;
624                            vec![PhysicalSortExpr {
625                                expr: time_index,
626                                options: SortOptions {
627                                    descending: false,
628                                    nulls_first: false,
629                                },
630                            }]
631                        };
632                        let arg = self.create_physical_expr_list(
633                            false,
634                            &aggr.params.args,
635                            input_dfschema,
636                            session,
637                            planning_ctx,
638                        )?;
639                        // first_value/last_value has only one param.
640                        // The param have been checked by datafusion in logical plan stage.
641                        // We can safely assume that there is only one element here.
642                        AggregateExprBuilder::new(aggr.func.clone(), arg)
643                            .schema(input_schema.clone())
644                            .order_by(order_by)
645                            .alias(name)
646                            .build()
647                    }
648                    Some(aggr) => {
649                        let order_by = if !aggr.params.order_by.is_empty() {
650                            aggr.params
651                                .order_by
652                                .iter()
653                                .map(|x| {
654                                    create_physical_sort_expr(
655                                        x,
656                                        input_dfschema.as_ref(),
657                                        session.execution_props(),
658                                        planning_ctx,
659                                    )
660                                })
661                                .collect::<DfResult<Vec<_>>>()?
662                        } else {
663                            vec![]
664                        };
665                        let distinct = aggr.params.distinct;
666                        // TODO(discord9): add default null treatment?
667
668                        let input_phy_exprs = self.create_physical_expr_list(
669                            aggr.func.name() == "count",
670                            &aggr.params.args,
671                            input_dfschema,
672                            session,
673                            planning_ctx,
674                        )?;
675                        AggregateExprBuilder::new(aggr.func.clone(), input_phy_exprs)
676                            .schema(input_schema.clone())
677                            .order_by(order_by)
678                            .with_distinct(distinct)
679                            .alias(name)
680                            .build()
681                    }
682                    None => Err(DataFusionError::Plan(format!(
683                        "Unexpected Expr: {} in RangeSelect",
684                        range_fn.expr
685                    ))),
686                }?;
687                Ok(RangeFnExec {
688                    expr: Arc::new(expr),
689                    range: range_fn.range.as_millis() as Millisecond,
690                    fill: range_fn.fill.clone(),
691                    need_cast: if range_fn.need_cast {
692                        Some(range_fn.data_type.clone())
693                    } else {
694                        None
695                    },
696                })
697            })
698            .collect::<DfResult<Vec<_>>>()?;
699        let schema_before_project = Arc::new(Schema::new(fields));
700        let schema = if let Some(project) = &self.schema_project {
701            Arc::new(schema_before_project.project(project)?)
702        } else {
703            schema_before_project.clone()
704        };
705        let by =
706            self.create_physical_expr_list(false, &self.by, input_dfschema, session, planning_ctx)?;
707        let cache = Arc::new(PlanProperties::new(
708            EquivalenceProperties::new(schema.clone()),
709            Partitioning::UnknownPartitioning(1),
710            EmissionType::Incremental,
711            Boundedness::Bounded,
712        ));
713        Ok(Arc::new(RangeSelectExec {
714            input: exec_input,
715            range_exec,
716            align: self.align.as_millis() as Millisecond,
717            align_to: self.align_to,
718            by,
719            time_index: self.time_index.clone(),
720            schema,
721            by_schema: Arc::new(Schema::new(by_fields)),
722            metric: ExecutionPlanMetricsSet::new(),
723            schema_before_project,
724            schema_project: self.schema_project.clone(),
725            cache,
726        }))
727    }
728}
729
730/// Range function expression.
731#[derive(Debug, Clone)]
732struct RangeFnExec {
733    expr: Arc<AggregateFunctionExpr>,
734    range: Millisecond,
735    fill: Option<Fill>,
736    need_cast: Option<DataType>,
737}
738
739impl RangeFnExec {
740    /// Returns the expressions to pass to the aggregator.
741    /// It also adds the order by expressions to the list of expressions.
742    /// Order-sensitive aggregators, such as `FIRST_VALUE(x ORDER BY y)` requires this.
743    fn expressions(&self) -> Vec<Arc<dyn PhysicalExpr>> {
744        let mut exprs = self.expr.expressions();
745        exprs.extend(self.expr.order_bys().iter().map(|sort| sort.expr.clone()));
746        exprs
747    }
748}
749
750impl Display for RangeFnExec {
751    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
752        if let Some(fill) = &self.fill {
753            write!(
754                f,
755                "{} RANGE {}s FILL {}",
756                self.expr.name(),
757                self.range / 1000,
758                fill
759            )
760        } else {
761            write!(f, "{} RANGE {}s", self.expr.name(), self.range / 1000)
762        }
763    }
764}
765
766#[derive(Debug)]
767pub struct RangeSelectExec {
768    input: Arc<dyn ExecutionPlan>,
769    range_exec: Vec<RangeFnExec>,
770    align: Millisecond,
771    align_to: i64,
772    time_index: String,
773    by: Vec<Arc<dyn PhysicalExpr>>,
774    schema: SchemaRef,
775    by_schema: SchemaRef,
776    metric: ExecutionPlanMetricsSet,
777    schema_project: Option<Vec<usize>>,
778    schema_before_project: SchemaRef,
779    cache: Arc<PlanProperties>,
780}
781
782impl DisplayAs for RangeSelectExec {
783    fn fmt_as(&self, t: DisplayFormatType, f: &mut std::fmt::Formatter) -> std::fmt::Result {
784        match t {
785            DisplayFormatType::Default
786            | DisplayFormatType::Verbose
787            | DisplayFormatType::TreeRender => {
788                write!(f, "RangeSelectExec: ")?;
789                let range_expr_strs: Vec<String> =
790                    self.range_exec.iter().map(RangeFnExec::to_string).collect();
791                let by: Vec<String> = self.by.iter().map(|e| e.to_string()).collect();
792                write!(
793                    f,
794                    "range_expr=[{}], align={}ms, align_to={}ms, align_by=[{}], time_index={}",
795                    range_expr_strs.join(", "),
796                    self.align,
797                    self.align_to,
798                    by.join(", "),
799                    self.time_index,
800                )?;
801            }
802        }
803        Ok(())
804    }
805}
806
807impl ExecutionPlan for RangeSelectExec {
808    fn schema(&self) -> SchemaRef {
809        self.schema.clone()
810    }
811
812    fn required_input_distribution(&self) -> Vec<Distribution> {
813        vec![Distribution::SinglePartition]
814    }
815
816    fn properties(&self) -> &Arc<PlanProperties> {
817        &self.cache
818    }
819
820    fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
821        vec![&self.input]
822    }
823
824    fn apply_expressions(
825        &self,
826        f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> DfResult<TreeNodeRecursion>,
827    ) -> DfResult<TreeNodeRecursion> {
828        apply_expression_roots(
829            self.range_exec
830                .iter()
831                .flat_map(RangeFnExec::expressions)
832                .chain(self.by.iter().cloned()),
833            f,
834        )
835    }
836
837    fn with_new_children(
838        self: Arc<Self>,
839        children: Vec<Arc<dyn ExecutionPlan>>,
840    ) -> datafusion_common::Result<Arc<dyn ExecutionPlan>> {
841        assert!(!children.is_empty());
842        Ok(Arc::new(Self {
843            input: children[0].clone(),
844            range_exec: self.range_exec.clone(),
845            time_index: self.time_index.clone(),
846            by: self.by.clone(),
847            align: self.align,
848            align_to: self.align_to,
849            schema: self.schema.clone(),
850            by_schema: self.by_schema.clone(),
851            metric: self.metric.clone(),
852            schema_before_project: self.schema_before_project.clone(),
853            schema_project: self.schema_project.clone(),
854            cache: self.cache.clone(),
855        }))
856    }
857
858    fn execute(
859        &self,
860        partition: usize,
861        context: Arc<TaskContext>,
862    ) -> DfResult<DfSendableRecordBatchStream> {
863        let baseline_metric = BaselineMetrics::new(&self.metric, partition);
864        let batch_size = context.session_config().batch_size();
865        let input = self.input.execute(partition, context)?;
866        let schema = input.schema();
867        let time_index = schema
868            .column_with_name(&self.time_index)
869            .ok_or(DataFusionError::Execution(
870                "time index column not found".into(),
871            ))?
872            .0;
873        let row_converter = RowConverter::new(
874            self.by_schema
875                .fields()
876                .iter()
877                .map(|f| SortField::new(f.data_type().clone()))
878                .collect(),
879        )?;
880        Ok(Box::pin(RangeSelectStream {
881            batch_size,
882            schema: self.schema.clone(),
883            range_exec: self.range_exec.clone(),
884            input,
885            random_state: RandomState::default(),
886            time_index,
887            align: self.align,
888            align_to: self.align_to,
889            by: self.by.clone(),
890            series_map: HashMap::new(),
891            exec_state: ExecutionState::ReadingInput,
892            num_not_null_rows: 0,
893            row_converter,
894            modify_map: HashMap::new(),
895            metric: baseline_metric,
896            schema_project: self.schema_project.clone(),
897            schema_before_project: self.schema_before_project.clone(),
898            output_batch: None,
899            output_batch_offset: 0,
900        }))
901    }
902
903    fn metrics(&self) -> Option<MetricsSet> {
904        Some(self.metric.clone_inner())
905    }
906
907    fn name(&self) -> &str {
908        "RanegSelectExec"
909    }
910}
911
912struct RangeSelectStream {
913    batch_size: usize,
914    /// the schema of output column
915    schema: SchemaRef,
916    range_exec: Vec<RangeFnExec>,
917    input: SendableRecordBatchStream,
918    /// Column index of TIME INDEX column's position in the input schema
919    time_index: usize,
920    /// the unit of `align` is millisecond
921    align: Millisecond,
922    align_to: i64,
923    by: Vec<Arc<dyn PhysicalExpr>>,
924    exec_state: ExecutionState,
925    /// Converter for the by values
926    row_converter: RowConverter,
927    random_state: RandomState,
928    /// key: time series's hash value
929    /// value: time series's state on different align_ts
930    series_map: HashMap<u64, SeriesState>,
931    /// key: `(hash of by rows, align_ts)`
932    /// value: `[row_ids]`
933    /// It is used to record the data that needs to be aggregated in each time slot during the data update process
934    modify_map: HashMap<(u64, Millisecond), Vec<u32>>,
935    /// The number of rows of not null rows in the final output
936    num_not_null_rows: usize,
937    metric: BaselineMetrics,
938    schema_project: Option<Vec<usize>>,
939    schema_before_project: SchemaRef,
940    output_batch: Option<RecordBatch>,
941    output_batch_offset: usize,
942}
943
944#[derive(Debug)]
945struct SeriesState {
946    /// by values written by `RowWriter`
947    row: OwnedRow,
948    /// key: align_ts
949    /// value: a vector, each element is a range_fn follow the order of `range_exec`
950    align_ts_accumulator: BTreeMap<Millisecond, Vec<Box<dyn Accumulator>>>,
951}
952
953/// Use `align_to` as time origin.
954/// According to `align` as time interval, produces aligned time.
955/// Combining the parameters related to the range query,
956/// determine for each `Accumulator` `(hash, align_ts)` define,
957/// which rows of data will be applied to it.
958fn produce_align_time(
959    align_to: i64,
960    range: Millisecond,
961    align: Millisecond,
962    ts_column: &TimestampMillisecondArray,
963    by_columns_hash: &[u64],
964    modify_map: &mut HashMap<(u64, Millisecond), Vec<u32>>,
965) {
966    modify_map.clear();
967    // make modify_map for range_fn[i]
968    for (row, hash) in by_columns_hash.iter().enumerate() {
969        let ts = ts_column.value(row);
970        let diff = ts - align_to;
971        // `div_euclid` equals `div_floor` for positive divisors (`align`).
972        let ith_slot = diff.div_euclid(align);
973        let mut align_ts = ith_slot * align + align_to;
974        while align_ts <= ts && ts < align_ts + range {
975            modify_map
976                .entry((*hash, align_ts))
977                .or_default()
978                .push(row as u32);
979            align_ts -= align;
980        }
981    }
982}
983
984fn cast_scalar_values(values: &mut [ScalarValue], data_type: &DataType) -> DfResult<()> {
985    let array = ScalarValue::iter_to_array(values.to_vec())?;
986    let cast_array = cast_with_options(&array, data_type, &CastOptions::default())?;
987    for (i, value) in values.iter_mut().enumerate() {
988        *value = ScalarValue::try_from_array(&cast_array, i)?;
989    }
990    Ok(())
991}
992
993impl RangeSelectStream {
994    fn evaluate_many(
995        &self,
996        batch: &RecordBatch,
997        exprs: &[Arc<dyn PhysicalExpr>],
998    ) -> DfResult<Vec<ArrayRef>> {
999        exprs
1000            .iter()
1001            .map(|expr| {
1002                let value = expr.evaluate(batch)?;
1003                value.into_array(batch.num_rows())
1004            })
1005            .collect::<DfResult<Vec<_>>>()
1006    }
1007
1008    fn update_range_context(&mut self, batch: RecordBatch) -> DfResult<()> {
1009        let _timer = self.metric.elapsed_compute().timer();
1010        let num_rows = batch.num_rows();
1011        let by_arrays = self.evaluate_many(&batch, &self.by)?;
1012        let mut hashes = vec![0; num_rows];
1013        create_hashes(&by_arrays, &self.random_state, &mut hashes)?;
1014        let by_rows = self.row_converter.convert_columns(&by_arrays)?;
1015        let mut ts_column = batch.column(self.time_index).clone();
1016        if !matches!(
1017            ts_column.data_type(),
1018            DataType::Timestamp(TimeUnit::Millisecond, _)
1019        ) {
1020            ts_column = compute::cast(
1021                ts_column.as_ref(),
1022                &DataType::Timestamp(TimeUnit::Millisecond, None),
1023            )?;
1024        }
1025        let ts_column_ref = ts_column
1026            .as_any()
1027            .downcast_ref::<TimestampMillisecondArray>()
1028            .ok_or_else(|| {
1029                DataFusionError::Execution(
1030                    "Time index Column downcast to TimestampMillisecondArray failed".into(),
1031                )
1032            })?;
1033        for i in 0..self.range_exec.len() {
1034            let args = self.evaluate_many(&batch, &self.range_exec[i].expressions())?;
1035            // use self.modify_map record (hash, align_ts) => [row_nums]
1036            produce_align_time(
1037                self.align_to,
1038                self.range_exec[i].range,
1039                self.align,
1040                ts_column_ref,
1041                &hashes,
1042                &mut self.modify_map,
1043            );
1044            // build modify_rows/modify_index/offsets for batch update
1045            let mut modify_rows = UInt32Builder::with_capacity(0);
1046            // (hash, align_ts, row_num)
1047            // row_num use to find a by value
1048            // So we just need to record the row_num of a modify row randomly, because they all have the same by value
1049            let mut modify_index = Vec::with_capacity(self.modify_map.len());
1050            let mut offsets = vec![0];
1051            let mut offset_so_far = 0;
1052            for ((hash, ts), modify) in &self.modify_map {
1053                modify_rows.append_slice(modify);
1054                offset_so_far += modify.len();
1055                offsets.push(offset_so_far);
1056                modify_index.push((*hash, *ts, modify[0]));
1057            }
1058            let modify_rows = modify_rows.finish();
1059            let args = take_arrays(&args, &modify_rows, None)?;
1060            modify_index.iter().zip(offsets.windows(2)).try_for_each(
1061                |((hash, ts, row), offset)| {
1062                    let (offset, length) = (offset[0], offset[1] - offset[0]);
1063                    let sliced_arrays: Vec<ArrayRef> = args
1064                        .iter()
1065                        .map(|array| array.slice(offset, length))
1066                        .collect();
1067                    let accumulators_map =
1068                        self.series_map.entry(*hash).or_insert_with(|| SeriesState {
1069                            row: by_rows.row(*row as usize).owned(),
1070                            align_ts_accumulator: BTreeMap::new(),
1071                        });
1072                    match accumulators_map.align_ts_accumulator.entry(*ts) {
1073                        Entry::Occupied(mut e) => {
1074                            let accumulators = e.get_mut();
1075                            accumulators[i].update_batch(&sliced_arrays)
1076                        }
1077                        Entry::Vacant(e) => {
1078                            self.num_not_null_rows += 1;
1079                            let mut accumulators = self
1080                                .range_exec
1081                                .iter()
1082                                .map(|range| range.expr.create_accumulator())
1083                                .collect::<DfResult<Vec<_>>>()?;
1084                            let result = accumulators[i].update_batch(&sliced_arrays);
1085                            e.insert(accumulators);
1086                            result
1087                        }
1088                    }
1089                },
1090            )?;
1091        }
1092        Ok(())
1093    }
1094
1095    fn generate_output(&mut self) -> DfResult<RecordBatch> {
1096        let _timer = self.metric.elapsed_compute().timer();
1097        if self.series_map.is_empty() {
1098            return Ok(RecordBatch::new_empty(self.schema.clone()));
1099        }
1100        // 1 for time index column
1101        let mut columns: Vec<Arc<dyn Array>> =
1102            Vec::with_capacity(1 + self.range_exec.len() + self.by.len());
1103        let mut ts_builder = TimestampMillisecondBuilder::with_capacity(self.num_not_null_rows);
1104        let mut all_scalar =
1105            vec![Vec::with_capacity(self.num_not_null_rows); self.range_exec.len()];
1106        let mut by_rows = Vec::with_capacity(self.num_not_null_rows);
1107        let mut start_index = 0;
1108        // If any range expr need fill, we need fill both the missing align_ts and null value.
1109        let need_fill_output = self.range_exec.iter().any(|range| range.fill.is_some());
1110        // The padding value for each accumulator
1111        let padding_values = self
1112            .range_exec
1113            .iter()
1114            .map(|e| e.expr.create_accumulator()?.evaluate())
1115            .collect::<DfResult<Vec<_>>>()?;
1116        for SeriesState {
1117            row,
1118            align_ts_accumulator,
1119        } in self.series_map.values_mut()
1120        {
1121            // skip empty time series
1122            if align_ts_accumulator.is_empty() {
1123                continue;
1124            }
1125            // find the first and last align_ts
1126            let begin_ts = *align_ts_accumulator.first_key_value().unwrap().0;
1127            let end_ts = *align_ts_accumulator.last_key_value().unwrap().0;
1128            let align_ts = if need_fill_output {
1129                // we need to fill empty align_ts which not data in that solt
1130                (begin_ts..=end_ts).step_by(self.align as usize).collect()
1131            } else {
1132                align_ts_accumulator.keys().copied().collect::<Vec<_>>()
1133            };
1134            for ts in &align_ts {
1135                if let Some(slot) = align_ts_accumulator.get_mut(ts) {
1136                    for (column, acc) in all_scalar.iter_mut().zip(slot.iter_mut()) {
1137                        column.push(acc.evaluate()?);
1138                    }
1139                } else {
1140                    // fill null in empty time solt
1141                    for (column, padding) in all_scalar.iter_mut().zip(padding_values.iter()) {
1142                        column.push(padding.clone())
1143                    }
1144                }
1145            }
1146            ts_builder.append_slice(&align_ts);
1147            // apply fill strategy on time series
1148            for (
1149                i,
1150                RangeFnExec {
1151                    fill, need_cast, ..
1152                },
1153            ) in self.range_exec.iter().enumerate()
1154            {
1155                let time_series_data =
1156                    &mut all_scalar[i][start_index..start_index + align_ts.len()];
1157                if let Some(data_type) = need_cast {
1158                    cast_scalar_values(time_series_data, data_type)?;
1159                }
1160                if let Some(fill) = fill {
1161                    fill.apply_fill_strategy(&align_ts, time_series_data)?;
1162                }
1163            }
1164            by_rows.resize(by_rows.len() + align_ts.len(), row.row());
1165            start_index += align_ts.len();
1166        }
1167        for column_scalar in all_scalar {
1168            columns.push(ScalarValue::iter_to_array(column_scalar)?);
1169        }
1170        let ts_column = ts_builder.finish();
1171        // output schema before project follow the order of range expr | time index | by columns
1172        let ts_column = compute::cast(
1173            &ts_column,
1174            self.schema_before_project.field(columns.len()).data_type(),
1175        )?;
1176        columns.push(ts_column);
1177        // RowConverter decodes dictionary sort fields to their value arrays. Re-encode them so
1178        // the physical batch continues to match the logical output schema.
1179        for by_column in self.row_converter.convert_rows(by_rows)? {
1180            let output_type = self.schema_before_project.field(columns.len()).data_type();
1181            if by_column.data_type() == output_type {
1182                columns.push(by_column);
1183            } else {
1184                columns.push(compute::cast(by_column.as_ref(), output_type)?);
1185            }
1186        }
1187        let output = RecordBatch::try_new(self.schema_before_project.clone(), columns)?;
1188        let project_output = if let Some(project) = &self.schema_project {
1189            output.project(project)?
1190        } else {
1191            output
1192        };
1193        Ok(project_output)
1194    }
1195
1196    fn next_output_batch(&mut self) -> DfResult<Option<RecordBatch>> {
1197        if self.output_batch.is_none() {
1198            self.output_batch = Some(self.generate_output()?);
1199            self.output_batch_offset = 0;
1200        }
1201
1202        let num_rows = self.output_batch.as_ref().unwrap().num_rows();
1203        if num_rows == 0 {
1204            self.output_batch = None;
1205            self.output_batch_offset = 0;
1206            return Ok(None);
1207        }
1208
1209        if self.output_batch_offset == 0 && num_rows <= self.batch_size {
1210            return Ok(self.output_batch.take());
1211        }
1212
1213        let offset = self.output_batch_offset;
1214        let len = (num_rows - offset).min(self.batch_size);
1215        let batch = self.output_batch.as_ref().unwrap().slice(offset, len);
1216        self.output_batch_offset += len;
1217
1218        if self.output_batch_offset >= num_rows {
1219            self.output_batch = None;
1220            self.output_batch_offset = 0;
1221        }
1222
1223        Ok(Some(batch))
1224    }
1225}
1226
1227enum ExecutionState {
1228    ReadingInput,
1229    ProducingOutput,
1230    Done,
1231}
1232
1233impl RecordBatchStream for RangeSelectStream {
1234    fn schema(&self) -> SchemaRef {
1235        self.schema.clone()
1236    }
1237}
1238
1239impl Stream for RangeSelectStream {
1240    type Item = DataFusionResult<RecordBatch>;
1241
1242    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
1243        loop {
1244            match self.exec_state {
1245                ExecutionState::ReadingInput => {
1246                    match ready!(self.input.poll_next_unpin(cx)) {
1247                        // new batch to aggregate
1248                        Some(Ok(batch)) => {
1249                            if let Err(e) = self.update_range_context(batch) {
1250                                common_telemetry::debug!(
1251                                    "RangeSelectStream cannot update range context, schema: {:?}, err: {:?}",
1252                                    self.schema,
1253                                    e
1254                                );
1255                                return Poll::Ready(Some(Err(e)));
1256                            }
1257                        }
1258                        // inner had error, return to caller
1259                        Some(Err(e)) => return Poll::Ready(Some(Err(e))),
1260                        // inner is done, producing output
1261                        None => {
1262                            self.exec_state = ExecutionState::ProducingOutput;
1263                        }
1264                    }
1265                }
1266                ExecutionState::ProducingOutput => {
1267                    let result = self.next_output_batch();
1268                    return match result {
1269                        // made output
1270                        Ok(Some(batch)) => {
1271                            if self.output_batch.is_none() {
1272                                self.exec_state = ExecutionState::Done;
1273                            }
1274                            Poll::Ready(Some(Ok(batch)))
1275                        }
1276                        Ok(None) => {
1277                            self.exec_state = ExecutionState::Done;
1278                            Poll::Ready(None)
1279                        }
1280                        // error making output
1281                        Err(error) => Poll::Ready(Some(Err(error))),
1282                    };
1283                }
1284                ExecutionState::Done => return Poll::Ready(None),
1285            }
1286        }
1287    }
1288}
1289
1290#[cfg(test)]
1291mod test {
1292    macro_rules! nullable_array {
1293        ($builder:ident,) => {
1294        };
1295        ($array_type:ident ; $($tail:tt)*) => {
1296            paste::item! {
1297                {
1298                    let mut builder = arrow::array::[<$array_type Builder>]::new();
1299                    nullable_array!(builder, $($tail)*);
1300                    builder.finish()
1301                }
1302            }
1303        };
1304        ($builder:ident, null) => {
1305            $builder.append_null();
1306        };
1307        ($builder:ident, null, $($tail:tt)*) => {
1308            $builder.append_null();
1309            nullable_array!($builder, $($tail)*);
1310        };
1311        ($builder:ident, $value:literal) => {
1312            $builder.append_value($value);
1313        };
1314        ($builder:ident, $value:literal, $($tail:tt)*) => {
1315            $builder.append_value($value);
1316            nullable_array!($builder, $($tail)*);
1317        };
1318    }
1319
1320    use std::sync::Arc;
1321
1322    use arrow_schema::SortOptions;
1323    use datafusion::arrow::datatypes::{
1324        ArrowPrimitiveType, DataType, Field, Schema, TimestampMillisecondType,
1325    };
1326    use datafusion::datasource::memory::MemorySourceConfig;
1327    use datafusion::datasource::source::DataSourceExec;
1328    use datafusion::functions_aggregate::{first_last, min_max};
1329    use datafusion::physical_plan::sorts::sort::SortExec;
1330    use datafusion::prelude::SessionContext;
1331    use datafusion_physical_expr::PhysicalSortExpr;
1332    use datafusion_physical_expr::expressions::Column;
1333    use datatypes::arrow::array::{Float64Array, Int64Array, TimestampMillisecondArray};
1334    use datatypes::arrow_array::StringArray;
1335
1336    use super::*;
1337
1338    const TIME_INDEX_COLUMN: &str = "timestamp";
1339
1340    fn prepare_test_data(is_float: bool, is_gap: bool) -> DataSourceExec {
1341        let schema = Arc::new(Schema::new(vec![
1342            Field::new(TIME_INDEX_COLUMN, TimestampMillisecondType::DATA_TYPE, true),
1343            Field::new(
1344                "value",
1345                if is_float {
1346                    DataType::Float64
1347                } else {
1348                    DataType::Int64
1349                },
1350                true,
1351            ),
1352            Field::new("host", DataType::Utf8, true),
1353        ]));
1354        let timestamp_column: Arc<dyn Array> = if !is_gap {
1355            Arc::new(TimestampMillisecondArray::from(vec![
1356                0, 5_000, 10_000, 15_000, 20_000, // host 1 every 5s
1357                0, 5_000, 10_000, 15_000, 20_000, // host 2 every 5s
1358            ])) as _
1359        } else {
1360            Arc::new(TimestampMillisecondArray::from(vec![
1361                0, 15_000, // host 1 every 5s, missing data on 5_000, 10_000
1362                0, 15_000, // host 2 every 5s, missing data on 5_000, 10_000
1363            ])) as _
1364        };
1365        let mut host = vec!["host1"; timestamp_column.len() / 2];
1366        host.extend(vec!["host2"; timestamp_column.len() / 2]);
1367        let mut value_column: Arc<dyn Array> = if is_gap {
1368            Arc::new(nullable_array!(Int64;
1369                0, 6, // data for host 1
1370                6, 12 // data for host 2
1371            )) as _
1372        } else {
1373            Arc::new(nullable_array!(Int64;
1374                0, null, 1, null, 2, // data for host 1
1375                3, null, 4, null, 5 // data for host 2
1376            )) as _
1377        };
1378        if is_float {
1379            value_column =
1380                cast_with_options(&value_column, &DataType::Float64, &CastOptions::default())
1381                    .unwrap();
1382        }
1383        let host_column: Arc<dyn Array> = Arc::new(StringArray::from(host)) as _;
1384        let data = RecordBatch::try_new(
1385            schema.clone(),
1386            vec![timestamp_column, value_column, host_column],
1387        )
1388        .unwrap();
1389
1390        DataSourceExec::new(Arc::new(
1391            MemorySourceConfig::try_new(&[vec![data]], schema, None).unwrap(),
1392        ))
1393    }
1394
1395    fn prepare_empty_test_data(is_float: bool) -> DataSourceExec {
1396        let schema = Arc::new(Schema::new(vec![
1397            Field::new(TIME_INDEX_COLUMN, TimestampMillisecondType::DATA_TYPE, true),
1398            Field::new(
1399                "value",
1400                if is_float {
1401                    DataType::Float64
1402                } else {
1403                    DataType::Int64
1404                },
1405                true,
1406            ),
1407            Field::new("host", DataType::Utf8, true),
1408        ]));
1409        let timestamp_column: Arc<dyn Array> =
1410            Arc::new(TimestampMillisecondArray::from(Vec::<i64>::new())) as _;
1411        let value_column: Arc<dyn Array> = if is_float {
1412            Arc::new(Float64Array::from(Vec::<Option<f64>>::new())) as _
1413        } else {
1414            Arc::new(Int64Array::from(Vec::<Option<i64>>::new())) as _
1415        };
1416        let host_column: Arc<dyn Array> =
1417            Arc::new(StringArray::from(Vec::<Option<&str>>::new())) as _;
1418        let data = RecordBatch::try_new(
1419            schema.clone(),
1420            vec![timestamp_column, value_column, host_column],
1421        )
1422        .unwrap();
1423
1424        DataSourceExec::new(Arc::new(
1425            MemorySourceConfig::try_new(&[vec![data]], schema, None).unwrap(),
1426        ))
1427    }
1428
1429    async fn collect_range_select_test(
1430        range1: Millisecond,
1431        range2: Millisecond,
1432        align: Millisecond,
1433        fill: Option<Fill>,
1434        is_float: bool,
1435        is_gap: bool,
1436        batch_size: usize,
1437    ) -> Vec<RecordBatch> {
1438        let data_type = if is_float {
1439            DataType::Float64
1440        } else {
1441            DataType::Int64
1442        };
1443        let (need_cast, schema_data_type) = if !is_float && matches!(fill, Some(Fill::Linear)) {
1444            // data_type = DataType::Float64;
1445            (Some(DataType::Float64), DataType::Float64)
1446        } else {
1447            (None, data_type.clone())
1448        };
1449        let memory_exec = Arc::new(prepare_test_data(is_float, is_gap));
1450        let schema = Arc::new(Schema::new(vec![
1451            Field::new("MIN(value)", schema_data_type.clone(), true),
1452            Field::new("MAX(value)", schema_data_type, true),
1453            Field::new(TIME_INDEX_COLUMN, TimestampMillisecondType::DATA_TYPE, true),
1454            Field::new("host", DataType::Utf8, true),
1455        ]));
1456        let cache = Arc::new(PlanProperties::new(
1457            EquivalenceProperties::new(schema.clone()),
1458            Partitioning::UnknownPartitioning(1),
1459            EmissionType::Incremental,
1460            Boundedness::Bounded,
1461        ));
1462        let input_schema = memory_exec.schema().clone();
1463        let range_select_exec = Arc::new(RangeSelectExec {
1464            input: memory_exec,
1465            range_exec: vec![
1466                RangeFnExec {
1467                    expr: Arc::new(
1468                        AggregateExprBuilder::new(
1469                            min_max::min_udaf(),
1470                            vec![Arc::new(Column::new("value", 1))],
1471                        )
1472                        .schema(input_schema.clone())
1473                        .alias("MIN(value)")
1474                        .build()
1475                        .unwrap(),
1476                    ),
1477                    range: range1,
1478                    fill: fill.clone(),
1479                    need_cast: need_cast.clone(),
1480                },
1481                RangeFnExec {
1482                    expr: Arc::new(
1483                        AggregateExprBuilder::new(
1484                            min_max::max_udaf(),
1485                            vec![Arc::new(Column::new("value", 1))],
1486                        )
1487                        .schema(input_schema.clone())
1488                        .alias("MAX(value)")
1489                        .build()
1490                        .unwrap(),
1491                    ),
1492                    range: range2,
1493                    fill,
1494                    need_cast,
1495                },
1496            ],
1497            align,
1498            align_to: 0,
1499            by: vec![Arc::new(Column::new("host", 2))],
1500            time_index: TIME_INDEX_COLUMN.to_string(),
1501            schema: schema.clone(),
1502            schema_before_project: schema.clone(),
1503            schema_project: None,
1504            by_schema: Arc::new(Schema::new(vec![Field::new("host", DataType::Utf8, true)])),
1505            metric: ExecutionPlanMetricsSet::new(),
1506            cache,
1507        });
1508        let sort_exec = SortExec::new(
1509            [
1510                PhysicalSortExpr {
1511                    expr: Arc::new(Column::new("host", 3)),
1512                    options: SortOptions {
1513                        descending: false,
1514                        nulls_first: true,
1515                    },
1516                },
1517                PhysicalSortExpr {
1518                    expr: Arc::new(Column::new(TIME_INDEX_COLUMN, 2)),
1519                    options: SortOptions {
1520                        descending: false,
1521                        nulls_first: true,
1522                    },
1523                },
1524            ]
1525            .into(),
1526            range_select_exec,
1527        );
1528        let session_context = SessionContext::new_with_config(
1529            datafusion::execution::config::SessionConfig::new().with_batch_size(batch_size),
1530        );
1531        datafusion::physical_plan::collect(Arc::new(sort_exec), session_context.task_ctx())
1532            .await
1533            .unwrap()
1534    }
1535
1536    async fn do_range_select_test(
1537        range1: Millisecond,
1538        range2: Millisecond,
1539        align: Millisecond,
1540        fill: Option<Fill>,
1541        is_float: bool,
1542        is_gap: bool,
1543        expected: String,
1544    ) {
1545        let result =
1546            collect_range_select_test(range1, range2, align, fill, is_float, is_gap, 8192).await;
1547
1548        let result_literal = arrow::util::pretty::pretty_format_batches(&result)
1549            .unwrap()
1550            .to_string();
1551
1552        assert_eq!(result_literal, expected);
1553    }
1554
1555    #[tokio::test]
1556    async fn range_10s_align_1000s() {
1557        let expected = String::from(
1558            "+------------+------------+---------------------+-------+\
1559            \n| MIN(value) | MAX(value) | timestamp           | host  |\
1560            \n+------------+------------+---------------------+-------+\
1561            \n| 0.0        | 0.0        | 1970-01-01T00:00:00 | host1 |\
1562            \n| 3.0        | 3.0        | 1970-01-01T00:00:00 | host2 |\
1563            \n+------------+------------+---------------------+-------+",
1564        );
1565        do_range_select_test(
1566            10_000,
1567            10_000,
1568            1_000_000,
1569            Some(Fill::Null),
1570            true,
1571            false,
1572            expected,
1573        )
1574        .await;
1575    }
1576
1577    #[tokio::test]
1578    async fn range_fill_null() {
1579        let expected = String::from(
1580            "+------------+------------+---------------------+-------+\
1581            \n| MIN(value) | MAX(value) | timestamp           | host  |\
1582            \n+------------+------------+---------------------+-------+\
1583            \n| 0.0        |            | 1969-12-31T23:59:55 | host1 |\
1584            \n| 0.0        | 0.0        | 1970-01-01T00:00:00 | host1 |\
1585            \n| 1.0        |            | 1970-01-01T00:00:05 | host1 |\
1586            \n| 1.0        | 1.0        | 1970-01-01T00:00:10 | host1 |\
1587            \n| 2.0        |            | 1970-01-01T00:00:15 | host1 |\
1588            \n| 2.0        | 2.0        | 1970-01-01T00:00:20 | host1 |\
1589            \n| 3.0        |            | 1969-12-31T23:59:55 | host2 |\
1590            \n| 3.0        | 3.0        | 1970-01-01T00:00:00 | host2 |\
1591            \n| 4.0        |            | 1970-01-01T00:00:05 | host2 |\
1592            \n| 4.0        | 4.0        | 1970-01-01T00:00:10 | host2 |\
1593            \n| 5.0        |            | 1970-01-01T00:00:15 | host2 |\
1594            \n| 5.0        | 5.0        | 1970-01-01T00:00:20 | host2 |\
1595            \n+------------+------------+---------------------+-------+",
1596        );
1597        do_range_select_test(
1598            10_000,
1599            5_000,
1600            5_000,
1601            Some(Fill::Null),
1602            true,
1603            false,
1604            expected,
1605        )
1606        .await;
1607    }
1608
1609    #[tokio::test]
1610    async fn range_fill_prev() {
1611        let expected = String::from(
1612            "+------------+------------+---------------------+-------+\
1613            \n| MIN(value) | MAX(value) | timestamp           | host  |\
1614            \n+------------+------------+---------------------+-------+\
1615            \n| 0.0        |            | 1969-12-31T23:59:55 | host1 |\
1616            \n| 0.0        | 0.0        | 1970-01-01T00:00:00 | host1 |\
1617            \n| 1.0        | 0.0        | 1970-01-01T00:00:05 | host1 |\
1618            \n| 1.0        | 1.0        | 1970-01-01T00:00:10 | host1 |\
1619            \n| 2.0        | 1.0        | 1970-01-01T00:00:15 | host1 |\
1620            \n| 2.0        | 2.0        | 1970-01-01T00:00:20 | host1 |\
1621            \n| 3.0        |            | 1969-12-31T23:59:55 | host2 |\
1622            \n| 3.0        | 3.0        | 1970-01-01T00:00:00 | host2 |\
1623            \n| 4.0        | 3.0        | 1970-01-01T00:00:05 | host2 |\
1624            \n| 4.0        | 4.0        | 1970-01-01T00:00:10 | host2 |\
1625            \n| 5.0        | 4.0        | 1970-01-01T00:00:15 | host2 |\
1626            \n| 5.0        | 5.0        | 1970-01-01T00:00:20 | host2 |\
1627            \n+------------+------------+---------------------+-------+",
1628        );
1629        do_range_select_test(
1630            10_000,
1631            5_000,
1632            5_000,
1633            Some(Fill::Prev),
1634            true,
1635            false,
1636            expected,
1637        )
1638        .await;
1639    }
1640
1641    #[tokio::test]
1642    async fn range_fill_linear() {
1643        let expected = String::from(
1644            "+------------+------------+---------------------+-------+\
1645            \n| MIN(value) | MAX(value) | timestamp           | host  |\
1646            \n+------------+------------+---------------------+-------+\
1647            \n| 0.0        | -0.5       | 1969-12-31T23:59:55 | host1 |\
1648            \n| 0.0        | 0.0        | 1970-01-01T00:00:00 | host1 |\
1649            \n| 1.0        | 0.5        | 1970-01-01T00:00:05 | host1 |\
1650            \n| 1.0        | 1.0        | 1970-01-01T00:00:10 | host1 |\
1651            \n| 2.0        | 1.5        | 1970-01-01T00:00:15 | host1 |\
1652            \n| 2.0        | 2.0        | 1970-01-01T00:00:20 | host1 |\
1653            \n| 3.0        | 2.5        | 1969-12-31T23:59:55 | host2 |\
1654            \n| 3.0        | 3.0        | 1970-01-01T00:00:00 | host2 |\
1655            \n| 4.0        | 3.5        | 1970-01-01T00:00:05 | host2 |\
1656            \n| 4.0        | 4.0        | 1970-01-01T00:00:10 | host2 |\
1657            \n| 5.0        | 4.5        | 1970-01-01T00:00:15 | host2 |\
1658            \n| 5.0        | 5.0        | 1970-01-01T00:00:20 | host2 |\
1659            \n+------------+------------+---------------------+-------+",
1660        );
1661        do_range_select_test(
1662            10_000,
1663            5_000,
1664            5_000,
1665            Some(Fill::Linear),
1666            true,
1667            false,
1668            expected,
1669        )
1670        .await;
1671    }
1672
1673    #[tokio::test]
1674    async fn range_fill_integer_linear() {
1675        let expected = String::from(
1676            "+------------+------------+---------------------+-------+\
1677            \n| MIN(value) | MAX(value) | timestamp           | host  |\
1678            \n+------------+------------+---------------------+-------+\
1679            \n| 0.0        | -0.5       | 1969-12-31T23:59:55 | host1 |\
1680            \n| 0.0        | 0.0        | 1970-01-01T00:00:00 | host1 |\
1681            \n| 1.0        | 0.5        | 1970-01-01T00:00:05 | host1 |\
1682            \n| 1.0        | 1.0        | 1970-01-01T00:00:10 | host1 |\
1683            \n| 2.0        | 1.5        | 1970-01-01T00:00:15 | host1 |\
1684            \n| 2.0        | 2.0        | 1970-01-01T00:00:20 | host1 |\
1685            \n| 3.0        | 2.5        | 1969-12-31T23:59:55 | host2 |\
1686            \n| 3.0        | 3.0        | 1970-01-01T00:00:00 | host2 |\
1687            \n| 4.0        | 3.5        | 1970-01-01T00:00:05 | host2 |\
1688            \n| 4.0        | 4.0        | 1970-01-01T00:00:10 | host2 |\
1689            \n| 5.0        | 4.5        | 1970-01-01T00:00:15 | host2 |\
1690            \n| 5.0        | 5.0        | 1970-01-01T00:00:20 | host2 |\
1691            \n+------------+------------+---------------------+-------+",
1692        );
1693        do_range_select_test(
1694            10_000,
1695            5_000,
1696            5_000,
1697            Some(Fill::Linear),
1698            false,
1699            false,
1700            expected,
1701        )
1702        .await;
1703    }
1704
1705    #[tokio::test]
1706    async fn range_fill_const() {
1707        let expected = String::from(
1708            "+------------+------------+---------------------+-------+\
1709            \n| MIN(value) | MAX(value) | timestamp           | host  |\
1710            \n+------------+------------+---------------------+-------+\
1711            \n| 0.0        | 6.6        | 1969-12-31T23:59:55 | host1 |\
1712            \n| 0.0        | 0.0        | 1970-01-01T00:00:00 | host1 |\
1713            \n| 1.0        | 6.6        | 1970-01-01T00:00:05 | host1 |\
1714            \n| 1.0        | 1.0        | 1970-01-01T00:00:10 | host1 |\
1715            \n| 2.0        | 6.6        | 1970-01-01T00:00:15 | host1 |\
1716            \n| 2.0        | 2.0        | 1970-01-01T00:00:20 | host1 |\
1717            \n| 3.0        | 6.6        | 1969-12-31T23:59:55 | host2 |\
1718            \n| 3.0        | 3.0        | 1970-01-01T00:00:00 | host2 |\
1719            \n| 4.0        | 6.6        | 1970-01-01T00:00:05 | host2 |\
1720            \n| 4.0        | 4.0        | 1970-01-01T00:00:10 | host2 |\
1721            \n| 5.0        | 6.6        | 1970-01-01T00:00:15 | host2 |\
1722            \n| 5.0        | 5.0        | 1970-01-01T00:00:20 | host2 |\
1723            \n+------------+------------+---------------------+-------+",
1724        );
1725        do_range_select_test(
1726            10_000,
1727            5_000,
1728            5_000,
1729            Some(Fill::Const(ScalarValue::Float64(Some(6.6)))),
1730            true,
1731            false,
1732            expected,
1733        )
1734        .await;
1735    }
1736
1737    #[tokio::test]
1738    async fn range_fill_gap() {
1739        let expected = String::from(
1740            "+------------+------------+---------------------+-------+\
1741            \n| MIN(value) | MAX(value) | timestamp           | host  |\
1742            \n+------------+------------+---------------------+-------+\
1743            \n| 0.0        | 0.0        | 1970-01-01T00:00:00 | host1 |\
1744            \n| 6.0        | 6.0        | 1970-01-01T00:00:15 | host1 |\
1745            \n| 6.0        | 6.0        | 1970-01-01T00:00:00 | host2 |\
1746            \n| 12.0       | 12.0       | 1970-01-01T00:00:15 | host2 |\
1747            \n+------------+------------+---------------------+-------+",
1748        );
1749        do_range_select_test(5_000, 5_000, 5_000, None, true, true, expected).await;
1750        let expected = String::from(
1751            "+------------+------------+---------------------+-------+\
1752            \n| MIN(value) | MAX(value) | timestamp           | host  |\
1753            \n+------------+------------+---------------------+-------+\
1754            \n| 0.0        | 0.0        | 1970-01-01T00:00:00 | host1 |\
1755            \n|            |            | 1970-01-01T00:00:05 | host1 |\
1756            \n|            |            | 1970-01-01T00:00:10 | host1 |\
1757            \n| 6.0        | 6.0        | 1970-01-01T00:00:15 | host1 |\
1758            \n| 6.0        | 6.0        | 1970-01-01T00:00:00 | host2 |\
1759            \n|            |            | 1970-01-01T00:00:05 | host2 |\
1760            \n|            |            | 1970-01-01T00:00:10 | host2 |\
1761            \n| 12.0       | 12.0       | 1970-01-01T00:00:15 | host2 |\
1762            \n+------------+------------+---------------------+-------+",
1763        );
1764        do_range_select_test(5_000, 5_000, 5_000, Some(Fill::Null), true, true, expected).await;
1765        let expected = String::from(
1766            "+------------+------------+---------------------+-------+\
1767            \n| MIN(value) | MAX(value) | timestamp           | host  |\
1768            \n+------------+------------+---------------------+-------+\
1769            \n| 0.0        | 0.0        | 1970-01-01T00:00:00 | host1 |\
1770            \n| 0.0        | 0.0        | 1970-01-01T00:00:05 | host1 |\
1771            \n| 0.0        | 0.0        | 1970-01-01T00:00:10 | host1 |\
1772            \n| 6.0        | 6.0        | 1970-01-01T00:00:15 | host1 |\
1773            \n| 6.0        | 6.0        | 1970-01-01T00:00:00 | host2 |\
1774            \n| 6.0        | 6.0        | 1970-01-01T00:00:05 | host2 |\
1775            \n| 6.0        | 6.0        | 1970-01-01T00:00:10 | host2 |\
1776            \n| 12.0       | 12.0       | 1970-01-01T00:00:15 | host2 |\
1777            \n+------------+------------+---------------------+-------+",
1778        );
1779        do_range_select_test(5_000, 5_000, 5_000, Some(Fill::Prev), true, true, expected).await;
1780        let expected = String::from(
1781            "+------------+------------+---------------------+-------+\
1782            \n| MIN(value) | MAX(value) | timestamp           | host  |\
1783            \n+------------+------------+---------------------+-------+\
1784            \n| 0.0        | 0.0        | 1970-01-01T00:00:00 | host1 |\
1785            \n| 2.0        | 2.0        | 1970-01-01T00:00:05 | host1 |\
1786            \n| 4.0        | 4.0        | 1970-01-01T00:00:10 | host1 |\
1787            \n| 6.0        | 6.0        | 1970-01-01T00:00:15 | host1 |\
1788            \n| 6.0        | 6.0        | 1970-01-01T00:00:00 | host2 |\
1789            \n| 8.0        | 8.0        | 1970-01-01T00:00:05 | host2 |\
1790            \n| 10.0       | 10.0       | 1970-01-01T00:00:10 | host2 |\
1791            \n| 12.0       | 12.0       | 1970-01-01T00:00:15 | host2 |\
1792            \n+------------+------------+---------------------+-------+",
1793        );
1794        do_range_select_test(
1795            5_000,
1796            5_000,
1797            5_000,
1798            Some(Fill::Linear),
1799            true,
1800            true,
1801            expected,
1802        )
1803        .await;
1804        let expected = String::from(
1805            "+------------+------------+---------------------+-------+\
1806            \n| MIN(value) | MAX(value) | timestamp           | host  |\
1807            \n+------------+------------+---------------------+-------+\
1808            \n| 0.0        | 0.0        | 1970-01-01T00:00:00 | host1 |\
1809            \n| 6.0        | 6.0        | 1970-01-01T00:00:05 | host1 |\
1810            \n| 6.0        | 6.0        | 1970-01-01T00:00:10 | host1 |\
1811            \n| 6.0        | 6.0        | 1970-01-01T00:00:15 | host1 |\
1812            \n| 6.0        | 6.0        | 1970-01-01T00:00:00 | host2 |\
1813            \n| 6.0        | 6.0        | 1970-01-01T00:00:05 | host2 |\
1814            \n| 6.0        | 6.0        | 1970-01-01T00:00:10 | host2 |\
1815            \n| 12.0       | 12.0       | 1970-01-01T00:00:15 | host2 |\
1816            \n+------------+------------+---------------------+-------+",
1817        );
1818        do_range_select_test(
1819            5_000,
1820            5_000,
1821            5_000,
1822            Some(Fill::Const(ScalarValue::Float64(Some(6.0)))),
1823            true,
1824            true,
1825            expected,
1826        )
1827        .await;
1828    }
1829
1830    #[test]
1831    fn range_select_apply_expressions_visits_owned_roots() {
1832        let input = Arc::new(prepare_test_data(true, false));
1833        let input_schema = input.schema().clone();
1834        let schema = Arc::new(Schema::new(vec![Field::new(
1835            "FIRST_VALUE(value)",
1836            DataType::Float64,
1837            true,
1838        )]));
1839        let range_select = RangeSelectExec {
1840            input,
1841            range_exec: vec![RangeFnExec {
1842                expr: Arc::new(
1843                    AggregateExprBuilder::new(
1844                        first_last::first_value_udaf(),
1845                        vec![Arc::new(Column::new("value", 1))],
1846                    )
1847                    .schema(input_schema)
1848                    .order_by(vec![PhysicalSortExpr {
1849                        expr: Arc::new(Column::new(TIME_INDEX_COLUMN, 0)),
1850                        options: SortOptions::default(),
1851                    }])
1852                    .alias("FIRST_VALUE(value)")
1853                    .build()
1854                    .unwrap(),
1855                ),
1856                range: 10_000,
1857                fill: None,
1858                need_cast: None,
1859            }],
1860            align: 5_000,
1861            align_to: 0,
1862            time_index: TIME_INDEX_COLUMN.to_string(),
1863            by: vec![Arc::new(Column::new("host", 2))],
1864            schema: schema.clone(),
1865            by_schema: Arc::new(Schema::empty()),
1866            metric: ExecutionPlanMetricsSet::new(),
1867            schema_project: None,
1868            schema_before_project: schema.clone(),
1869            cache: Arc::new(PlanProperties::new(
1870                EquivalenceProperties::new(schema),
1871                Partitioning::UnknownPartitioning(1),
1872                EmissionType::Incremental,
1873                Boundedness::Bounded,
1874            )),
1875        };
1876        assert_eq!(range_select.range_exec[0].expr.order_bys().len(), 1);
1877
1878        let mut visited = Vec::new();
1879        assert_eq!(
1880            range_select
1881                .apply_expressions(&mut |expr| {
1882                    visited.push(expr.to_string());
1883                    Ok(TreeNodeRecursion::Continue)
1884                })
1885                .unwrap(),
1886            TreeNodeRecursion::Continue
1887        );
1888        assert_eq!(visited, ["value@1", "timestamp@0", "host@2"]);
1889
1890        let mut stopped = Vec::new();
1891        assert_eq!(
1892            range_select
1893                .apply_expressions(&mut |expr| {
1894                    stopped.push(expr.to_string());
1895                    Ok(TreeNodeRecursion::Stop)
1896                })
1897                .unwrap(),
1898            TreeNodeRecursion::Stop
1899        );
1900        assert_eq!(stopped, ["value@1"]);
1901
1902        assert_eq!(
1903            range_select
1904                .apply_expressions(&mut |_| {
1905                    Err(DataFusionError::Execution("apply failure".into()))
1906                })
1907                .unwrap_err()
1908                .to_string(),
1909            "Execution error: apply failure"
1910        );
1911    }
1912
1913    #[tokio::test]
1914    async fn range_select_respects_session_batch_size() {
1915        let result =
1916            collect_range_select_test(10_000, 5_000, 5_000, Some(Fill::Null), true, false, 3).await;
1917
1918        let row_counts = result
1919            .iter()
1920            .map(|batch| batch.num_rows())
1921            .collect::<Vec<_>>();
1922        assert_eq!(vec![3, 3, 3, 3], row_counts);
1923    }
1924
1925    #[tokio::test]
1926    async fn range_select_skips_empty_output_batch() {
1927        let memory_exec = Arc::new(prepare_empty_test_data(true));
1928        let schema = Arc::new(Schema::new(vec![
1929            Field::new("MIN(value)", DataType::Float64, true),
1930            Field::new("MAX(value)", DataType::Float64, true),
1931            Field::new(TIME_INDEX_COLUMN, TimestampMillisecondType::DATA_TYPE, true),
1932            Field::new("host", DataType::Utf8, true),
1933        ]));
1934        let cache = Arc::new(PlanProperties::new(
1935            EquivalenceProperties::new(schema.clone()),
1936            Partitioning::UnknownPartitioning(1),
1937            EmissionType::Incremental,
1938            Boundedness::Bounded,
1939        ));
1940        let input_schema = memory_exec.schema().clone();
1941        let range_select_exec = Arc::new(RangeSelectExec {
1942            input: memory_exec,
1943            range_exec: vec![
1944                RangeFnExec {
1945                    expr: Arc::new(
1946                        AggregateExprBuilder::new(
1947                            min_max::min_udaf(),
1948                            vec![Arc::new(Column::new("value", 1))],
1949                        )
1950                        .schema(input_schema.clone())
1951                        .alias("MIN(value)")
1952                        .build()
1953                        .unwrap(),
1954                    ),
1955                    range: 10_000,
1956                    fill: Some(Fill::Null),
1957                    need_cast: None,
1958                },
1959                RangeFnExec {
1960                    expr: Arc::new(
1961                        AggregateExprBuilder::new(
1962                            min_max::max_udaf(),
1963                            vec![Arc::new(Column::new("value", 1))],
1964                        )
1965                        .schema(input_schema)
1966                        .alias("MAX(value)")
1967                        .build()
1968                        .unwrap(),
1969                    ),
1970                    range: 5_000,
1971                    fill: Some(Fill::Null),
1972                    need_cast: None,
1973                },
1974            ],
1975            align: 5_000,
1976            align_to: 0,
1977            by: vec![Arc::new(Column::new("host", 2))],
1978            time_index: TIME_INDEX_COLUMN.to_string(),
1979            schema: schema.clone(),
1980            schema_before_project: schema.clone(),
1981            schema_project: None,
1982            by_schema: Arc::new(Schema::new(vec![Field::new("host", DataType::Utf8, true)])),
1983            metric: ExecutionPlanMetricsSet::new(),
1984            cache,
1985        });
1986        let session_context = SessionContext::new();
1987        let result =
1988            datafusion::physical_plan::collect(range_select_exec, session_context.task_ctx())
1989                .await
1990                .unwrap();
1991
1992        assert!(result.is_empty());
1993    }
1994
1995    #[test]
1996    fn fill_test() {
1997        assert!(Fill::try_from_str("", &DataType::UInt8).unwrap().is_none());
1998        assert!(Fill::try_from_str("Linear", &DataType::UInt8).unwrap() == Some(Fill::Linear));
1999        assert_eq!(
2000            Fill::try_from_str("Linear", &DataType::Boolean)
2001                .unwrap_err()
2002                .to_string(),
2003            "Error during planning: Use FILL LINEAR on Non-numeric DataType Boolean"
2004        );
2005        assert_eq!(
2006            Fill::try_from_str("WHAT", &DataType::UInt8)
2007                .unwrap_err()
2008                .to_string(),
2009            "Error during planning: WHAT is not a valid fill option, fail to convert to a const value. { Arrow error: Cast error: Cannot cast string 'WHAT' to value of UInt8 type }"
2010        );
2011        assert_eq!(
2012            Fill::try_from_str("8.0", &DataType::UInt8)
2013                .unwrap_err()
2014                .to_string(),
2015            "Error during planning: 8.0 is not a valid fill option, fail to convert to a const value. { Arrow error: Cast error: Cannot cast string '8.0' to value of UInt8 type }"
2016        );
2017        assert!(
2018            Fill::try_from_str("8", &DataType::UInt8).unwrap()
2019                == Some(Fill::Const(ScalarValue::UInt8(Some(8))))
2020        );
2021        let mut test1 = vec![
2022            ScalarValue::UInt8(Some(8)),
2023            ScalarValue::UInt8(None),
2024            ScalarValue::UInt8(Some(9)),
2025        ];
2026        Fill::Null.apply_fill_strategy(&[], &mut test1).unwrap();
2027        assert_eq!(test1[1], ScalarValue::UInt8(None));
2028        Fill::Prev.apply_fill_strategy(&[], &mut test1).unwrap();
2029        assert_eq!(test1[1], ScalarValue::UInt8(Some(8)));
2030        test1[1] = ScalarValue::UInt8(None);
2031        Fill::Const(ScalarValue::UInt8(Some(10)))
2032            .apply_fill_strategy(&[], &mut test1)
2033            .unwrap();
2034        assert_eq!(test1[1], ScalarValue::UInt8(Some(10)));
2035    }
2036
2037    #[test]
2038    fn test_fill_linear() {
2039        let ts = vec![1, 2, 3, 4, 5];
2040        let mut test = vec![
2041            ScalarValue::Float32(Some(1.0)),
2042            ScalarValue::Float32(None),
2043            ScalarValue::Float32(Some(3.0)),
2044            ScalarValue::Float32(None),
2045            ScalarValue::Float32(Some(5.0)),
2046        ];
2047        Fill::Linear.apply_fill_strategy(&ts, &mut test).unwrap();
2048        let mut test1 = vec![
2049            ScalarValue::Float32(None),
2050            ScalarValue::Float32(Some(2.0)),
2051            ScalarValue::Float32(None),
2052            ScalarValue::Float32(Some(4.0)),
2053            ScalarValue::Float32(None),
2054        ];
2055        Fill::Linear.apply_fill_strategy(&ts, &mut test1).unwrap();
2056        assert_eq!(test, test1);
2057        // test linear interpolation on irregularly spaced ts/data
2058        let ts = vec![
2059            1,   // None
2060            3,   // 1.0
2061            8,   // 11.0
2062            30,  // None
2063            88,  // 10.0
2064            108, // 5.0
2065            128, // None
2066        ];
2067        let mut test = vec![
2068            ScalarValue::Float64(None),
2069            ScalarValue::Float64(Some(1.0)),
2070            ScalarValue::Float64(Some(11.0)),
2071            ScalarValue::Float64(None),
2072            ScalarValue::Float64(Some(10.0)),
2073            ScalarValue::Float64(Some(5.0)),
2074            ScalarValue::Float64(None),
2075        ];
2076        Fill::Linear.apply_fill_strategy(&ts, &mut test).unwrap();
2077        let data: Vec<_> = test
2078            .into_iter()
2079            .map(|x| {
2080                let ScalarValue::Float64(Some(f)) = x else {
2081                    unreachable!()
2082                };
2083                f
2084            })
2085            .collect();
2086        assert_eq!(data, vec![-3.0, 1.0, 11.0, 10.725, 10.0, 5.0, 0.0]);
2087        // test corner case
2088        let ts = vec![1];
2089        let test = vec![ScalarValue::Float32(None)];
2090        let mut test1 = test.clone();
2091        Fill::Linear.apply_fill_strategy(&ts, &mut test1).unwrap();
2092        assert_eq!(test, test1);
2093    }
2094}