1use 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 pub fn apply_fill_strategy(&self, ts: &[i64], data: &mut [ScalarValue]) -> DfResult<()> {
115 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 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 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 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 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 if let Some(end) = head {
181 linear_interpolation(ts, data, end, end + 1, 0, end)?;
182 }
183 if let Some(start) = tail {
185 linear_interpolation(ts, data, start - 2, start - 1, start, data.len())?;
186 }
187 Ok(())
188 }
189}
190
191fn 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 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 pub name: String,
233 pub data_type: DataType,
234 pub expr: Expr,
235 pub range: Duration,
236 pub fill: Option<Fill>,
237 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 pub input: Arc<LogicalPlan>,
278 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 pub schema_project: Option<Vec<usize>>,
291 pub schema_before_project: DFSchemaRef,
295}
296
297impl PartialOrd for RangeSelect {
298 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
299 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 !matches!(fill, Some(Fill::Const(..))),
373 );
374 Ok((None, Arc::new(field)))
375 },
376 )
377 .collect::<DfResult<Vec<_>>>()?;
378 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 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 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 #[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 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 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 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#[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 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 schema: SchemaRef,
916 range_exec: Vec<RangeFnExec>,
917 input: SendableRecordBatchStream,
918 time_index: usize,
920 align: Millisecond,
922 align_to: i64,
923 by: Vec<Arc<dyn PhysicalExpr>>,
924 exec_state: ExecutionState,
925 row_converter: RowConverter,
927 random_state: RandomState,
928 series_map: HashMap<u64, SeriesState>,
931 modify_map: HashMap<(u64, Millisecond), Vec<u32>>,
935 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 row: OwnedRow,
948 align_ts_accumulator: BTreeMap<Millisecond, Vec<Box<dyn Accumulator>>>,
951}
952
953fn 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 for (row, hash) in by_columns_hash.iter().enumerate() {
969 let ts = ts_column.value(row);
970 let diff = ts - align_to;
971 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 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 let mut modify_rows = UInt32Builder::with_capacity(0);
1046 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 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 let need_fill_output = self.range_exec.iter().any(|range| range.fill.is_some());
1110 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 if align_ts_accumulator.is_empty() {
1123 continue;
1124 }
1125 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 (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 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 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 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 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 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 Some(Err(e)) => return Poll::Ready(Some(Err(e))),
1260 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 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 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, 0, 5_000, 10_000, 15_000, 20_000, ])) as _
1359 } else {
1360 Arc::new(TimestampMillisecondArray::from(vec![
1361 0, 15_000, 0, 15_000, ])) 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, 6, 12 )) as _
1372 } else {
1373 Arc::new(nullable_array!(Int64;
1374 0, null, 1, null, 2, 3, null, 4, null, 5 )) 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 (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 let ts = vec![
2059 1, 3, 8, 30, 88, 108, 128, ];
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 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}