1use std::collections::HashSet;
16use std::pin::Pin;
17use std::sync::Arc;
18use std::task::{Context, Poll};
19
20use common_telemetry::{debug, warn};
21use datafusion::arrow::array::{Array, ArrayRef, Int64Array, TimestampMillisecondArray};
22use datafusion::arrow::compute;
23use datafusion::arrow::datatypes::{DataType, Field, SchemaRef, TimeUnit};
24use datafusion::arrow::error::ArrowError;
25use datafusion::arrow::record_batch::RecordBatch;
26use datafusion::common::stats::Precision;
27use datafusion::common::tree_node::TreeNodeRecursion;
28use datafusion::common::{DFSchema, DFSchemaRef, TableReference};
29use datafusion::error::{DataFusionError, Result as DataFusionResult};
30use datafusion::execution::context::TaskContext;
31use datafusion::logical_expr::{EmptyRelation, Expr, LogicalPlan, UserDefinedLogicalNodeCore};
32use datafusion::physical_expr::EquivalenceProperties;
33use datafusion::physical_plan::metrics::{
34 BaselineMetrics, Count, ExecutionPlanMetricsSet, MetricBuilder, MetricValue, MetricsSet,
35};
36use datafusion::physical_plan::{
37 ChildStats, DisplayAs, DisplayFormatType, Distribution, ExecutionPlan,
38 InputDistributionRequirements, PhysicalExpr, PlanProperties, RecordBatchStream,
39 SendableRecordBatchStream, Statistics, StatisticsArgs,
40};
41use datafusion_expr::ident;
42use datatypes::timestamp::timestamp_array_to_primitive;
43use futures::{Stream, StreamExt, ready};
44use greptime_proto::substrait_extension as pb;
45use prost::Message;
46use snafu::ResultExt;
47
48use crate::error::{DeserializeSnafu, Result};
49use crate::extension_plan::{
50 METRIC_NUM_SERIES, Millisecond, local_offset, nanoseconds_per_native_tick, resolve_column_name,
51 serialize_column_index, timestamp_unit,
52};
53use crate::metrics::PROMQL_SERIES_COUNT;
54use crate::range_array::RangeArray;
55
56#[derive(Debug, PartialEq, Eq, Hash)]
66pub struct RangeManipulate {
67 start: Millisecond,
68 end: Millisecond,
69 interval: Millisecond,
70 range: Millisecond,
71 offset: Millisecond,
72 time_index: String,
73 field_columns: Vec<String>,
74 input: LogicalPlan,
75 output_schema: DFSchemaRef,
76 unfix: Option<UnfixIndices>,
77}
78
79#[derive(Debug, Clone, PartialEq, Eq, Hash)]
80struct UnfixIndices {
81 pub time_index_idx: u64,
82 pub tag_column_indices: Vec<u64>,
83}
84
85impl RangeManipulate {
86 #[allow(clippy::too_many_arguments)]
87 pub fn new(
88 start: Millisecond,
89 end: Millisecond,
90 interval: Millisecond,
91 offset: Millisecond,
92 range: Millisecond,
93 time_index: String,
94 field_columns: Vec<String>,
95 input: LogicalPlan,
96 ) -> DataFusionResult<Self> {
97 let output_schema =
98 Self::calculate_output_schema(input.schema(), &time_index, &field_columns)?;
99 Ok(Self {
100 start,
101 end,
102 interval,
103 range,
104 offset,
105 time_index,
106 field_columns,
107 input,
108 output_schema,
109 unfix: None,
110 })
111 }
112
113 pub const fn name() -> &'static str {
114 "RangeManipulate"
115 }
116
117 pub fn build_timestamp_range_name(time_index: &str) -> String {
118 format!("{time_index}_range")
119 }
120
121 pub fn internal_range_end_col_name() -> String {
122 "__internal_range_end".to_string()
123 }
124
125 fn range_timestamp_name(&self) -> String {
126 Self::build_timestamp_range_name(&self.time_index)
127 }
128
129 fn calculate_output_schema(
130 input_schema: &DFSchemaRef,
131 time_index: &str,
132 field_columns: &[String],
133 ) -> DataFusionResult<DFSchemaRef> {
134 let columns = input_schema.fields();
135 let mut new_columns = Vec::with_capacity(columns.len() + 1);
136 for i in 0..columns.len() {
137 let x = input_schema.qualified_field(i);
138 new_columns.push((x.0.cloned(), x.1.clone()));
139 }
140
141 let Some(ts_col_index) = input_schema.index_of_column_by_name(None, time_index) else {
144 return Err(datafusion::common::field_not_found(
145 None::<TableReference>,
146 time_index,
147 input_schema.as_ref(),
148 ));
149 };
150 let ts_col_field = &columns[ts_col_index];
151 let output_time_field = Arc::new(
152 ts_col_field
153 .as_ref()
154 .clone()
155 .with_data_type(DataType::Timestamp(TimeUnit::Millisecond, None)),
156 );
157 new_columns[ts_col_index] = (
158 input_schema.qualified_field(ts_col_index).0.cloned(),
159 output_time_field.clone(),
160 );
161 let timestamp_range_field = Field::new(
162 Self::build_timestamp_range_name(time_index),
163 RangeArray::convert_field(output_time_field.as_ref())
164 .data_type()
165 .clone(),
166 ts_col_field.is_nullable(),
167 );
168 new_columns.push((None, Arc::new(timestamp_range_field)));
169
170 for name in field_columns {
172 let Some(index) = input_schema.index_of_column_by_name(None, name) else {
173 return Err(datafusion::common::field_not_found(
174 None::<TableReference>,
175 name,
176 input_schema.as_ref(),
177 ));
178 };
179 new_columns[index] = (None, Arc::new(RangeArray::convert_field(&columns[index])));
180 }
181
182 Ok(Arc::new(DFSchema::new_with_metadata(
183 new_columns,
184 input_schema.metadata().clone(),
185 )?))
186 }
187
188 pub fn to_execution_plan(&self, exec_input: Arc<dyn ExecutionPlan>) -> Arc<dyn ExecutionPlan> {
189 let output_schema: SchemaRef = self.output_schema.inner().clone();
190 let properties = exec_input.properties();
191 let properties = Arc::new(PlanProperties::new(
192 EquivalenceProperties::new(output_schema.clone()),
193 properties.partitioning.clone(),
194 properties.emission_type,
195 properties.boundedness,
196 ));
197 Arc::new(RangeManipulateExec {
198 offset: self.offset,
199 start: self.start,
200 end: self.end,
201 interval: self.interval,
202 range: self.range,
203 time_index_column: self.time_index.clone(),
204 time_range_column: self.range_timestamp_name(),
205 field_columns: self.field_columns.clone(),
206 input: exec_input,
207 output_schema,
208 metric: ExecutionPlanMetricsSet::new(),
209 properties,
210 })
211 }
212
213 pub fn serialize(&self) -> Vec<u8> {
214 let time_index_idx = serialize_column_index(self.input.schema(), &self.time_index);
215
216 let tag_column_indices = self
217 .field_columns
218 .iter()
219 .map(|name| serialize_column_index(self.input.schema(), name))
220 .collect::<Vec<u64>>();
221
222 pb::RangeManipulate {
223 start: self.start,
224 end: self.end,
225 interval: self.interval,
226 range: self.range,
227 time_index_idx,
228 tag_column_indices,
229 ..Default::default()
230 }
231 .encode_to_vec()
232 }
233
234 pub fn deserialize(bytes: &[u8]) -> Result<Self> {
235 let pb_range_manipulate = pb::RangeManipulate::decode(bytes).context(DeserializeSnafu)?;
236 let empty_schema = Arc::new(DFSchema::empty());
237 let placeholder_plan = LogicalPlan::EmptyRelation(EmptyRelation {
238 produce_one_row: false,
239 schema: empty_schema.clone(),
240 });
241
242 let unfix = UnfixIndices {
243 time_index_idx: pb_range_manipulate.time_index_idx,
244 tag_column_indices: pb_range_manipulate.tag_column_indices.clone(),
245 };
246 debug!("RangeManipulate deserialize unfix: {:?}", unfix);
247
248 Ok(Self {
253 start: pb_range_manipulate.start,
254 end: pb_range_manipulate.end,
255 interval: pb_range_manipulate.interval,
256 range: pb_range_manipulate.range,
257 offset: 0,
258 time_index: String::new(),
259 field_columns: Vec::new(),
260 input: placeholder_plan,
261 output_schema: empty_schema,
262 unfix: Some(unfix),
263 })
264 }
265}
266
267impl PartialOrd for RangeManipulate {
268 fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
269 match self.start.partial_cmp(&other.start) {
271 Some(core::cmp::Ordering::Equal) => {}
272 ord => return ord,
273 }
274 match self.end.partial_cmp(&other.end) {
275 Some(core::cmp::Ordering::Equal) => {}
276 ord => return ord,
277 }
278 match self.interval.partial_cmp(&other.interval) {
279 Some(core::cmp::Ordering::Equal) => {}
280 ord => return ord,
281 }
282 match self.range.partial_cmp(&other.range) {
283 Some(core::cmp::Ordering::Equal) => {}
284 ord => return ord,
285 }
286 match self.offset.partial_cmp(&other.offset) {
287 Some(core::cmp::Ordering::Equal) => {}
288 ord => return ord,
289 }
290 match self.time_index.partial_cmp(&other.time_index) {
291 Some(core::cmp::Ordering::Equal) => {}
292 ord => return ord,
293 }
294 match self.field_columns.partial_cmp(&other.field_columns) {
295 Some(core::cmp::Ordering::Equal) => {}
296 ord => return ord,
297 }
298 self.input.partial_cmp(&other.input)
299 }
300}
301
302impl UserDefinedLogicalNodeCore for RangeManipulate {
303 fn name(&self) -> &str {
304 Self::name()
305 }
306
307 fn inputs(&self) -> Vec<&LogicalPlan> {
308 vec![&self.input]
309 }
310
311 fn schema(&self) -> &DFSchemaRef {
312 &self.output_schema
313 }
314
315 fn expressions(&self) -> Vec<Expr> {
316 if self.unfix.is_some() {
317 return vec![];
318 }
319
320 let mut exprs = Vec::with_capacity(1 + self.field_columns.len());
321 exprs.push(ident(&self.time_index));
322 exprs.extend(self.field_columns.iter().map(ident));
323 exprs
324 }
325
326 fn necessary_children_exprs(&self, output_columns: &[usize]) -> Option<Vec<Vec<usize>>> {
327 if self.unfix.is_some() {
328 return None;
329 }
330
331 let input_schema = self.input.schema();
332 let input_len = input_schema.fields().len();
333 let time_index_idx = input_schema.index_of_column_by_name(None, &self.time_index)?;
334
335 if output_columns.is_empty() {
336 let indices = (0..input_len).collect::<Vec<_>>();
337 return Some(vec![indices]);
338 }
339
340 let mut required = Vec::with_capacity(output_columns.len() + 1 + self.field_columns.len());
341 required.push(time_index_idx);
342 for value_column in &self.field_columns {
343 required.push(input_schema.index_of_column_by_name(None, value_column)?);
344 }
345 for &idx in output_columns {
346 if idx < input_len {
347 required.push(idx);
348 } else if idx == input_len {
349 required.push(time_index_idx);
351 } else {
352 warn!(
353 "Output column index {} is out of bounds for input schema with length {}",
354 idx, input_len
355 );
356 return None;
357 }
358 }
359
360 required.sort_unstable();
361 required.dedup();
362 Some(vec![required])
363 }
364
365 fn fmt_for_explain(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
366 write!(
367 f,
368 "PromRangeManipulate: req range=[{}..{}], interval=[{}], eval range=[{}], time index=[{}], values={:?}",
369 self.start, self.end, self.interval, self.range, self.time_index, self.field_columns
370 )
371 }
372
373 fn with_exprs_and_inputs(
374 &self,
375 _exprs: Vec<Expr>,
376 mut inputs: Vec<LogicalPlan>,
377 ) -> DataFusionResult<Self> {
378 if inputs.len() != 1 {
379 return Err(DataFusionError::Internal(
380 "RangeManipulate should have at exact one input".to_string(),
381 ));
382 }
383
384 let input: LogicalPlan = inputs.pop().unwrap();
385 let input_schema = input.schema();
386
387 if let Some(unfix) = &self.unfix {
388 let time_index = resolve_column_name(
390 unfix.time_index_idx,
391 input_schema,
392 "RangeManipulate",
393 "time index",
394 )?;
395
396 let field_columns = unfix
397 .tag_column_indices
398 .iter()
399 .map(|idx| resolve_column_name(*idx, input_schema, "RangeManipulate", "tag"))
400 .collect::<DataFusionResult<Vec<String>>>()?;
401
402 let output_schema =
403 Self::calculate_output_schema(input_schema, &time_index, &field_columns)?;
404
405 Ok(Self {
406 start: self.start,
407 end: self.end,
408 interval: self.interval,
409 range: self.range,
410 offset: local_offset(&input, &time_index),
411 time_index,
412 field_columns,
413 input,
414 output_schema,
415 unfix: None,
416 })
417 } else {
418 let output_schema =
419 Self::calculate_output_schema(input_schema, &self.time_index, &self.field_columns)?;
420
421 Ok(Self {
422 start: self.start,
423 end: self.end,
424 interval: self.interval,
425 range: self.range,
426 offset: self.offset,
427 time_index: self.time_index.clone(),
428 field_columns: self.field_columns.clone(),
429 input,
430 output_schema,
431 unfix: None,
432 })
433 }
434 }
435}
436
437#[derive(Debug)]
438pub struct RangeManipulateExec {
439 offset: Millisecond,
440 start: Millisecond,
441 end: Millisecond,
442 interval: Millisecond,
443 range: Millisecond,
444 time_index_column: String,
445 time_range_column: String,
446 field_columns: Vec<String>,
447
448 input: Arc<dyn ExecutionPlan>,
449 output_schema: SchemaRef,
450 metric: ExecutionPlanMetricsSet,
451 properties: Arc<PlanProperties>,
452}
453
454impl ExecutionPlan for RangeManipulateExec {
455 fn apply_expressions(
456 &self,
457 _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> datafusion_common::Result<TreeNodeRecursion>,
458 ) -> DataFusionResult<TreeNodeRecursion> {
459 Ok(TreeNodeRecursion::Continue)
460 }
461
462 fn schema(&self) -> SchemaRef {
463 self.output_schema.clone()
464 }
465
466 fn properties(&self) -> &Arc<PlanProperties> {
467 &self.properties
468 }
469
470 fn maintains_input_order(&self) -> Vec<bool> {
471 vec![true; self.children().len()]
472 }
473
474 fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
475 vec![&self.input]
476 }
477
478 fn input_distribution_requirements(&self) -> InputDistributionRequirements {
479 let input_requirement = self
480 .input
481 .input_distribution_requirements()
482 .into_per_child();
483 if input_requirement.is_empty() {
484 InputDistributionRequirements::new(vec![Distribution::UnspecifiedDistribution])
487 } else {
488 InputDistributionRequirements::new(input_requirement)
489 }
490 }
491
492 fn with_new_children(
493 self: Arc<Self>,
494 children: Vec<Arc<dyn ExecutionPlan>>,
495 ) -> DataFusionResult<Arc<dyn ExecutionPlan>> {
496 assert!(!children.is_empty());
497 let exec_input = children[0].clone();
498 let properties = exec_input.properties();
499 let properties = Arc::new(PlanProperties::new(
500 EquivalenceProperties::new(self.output_schema.clone()),
501 properties.partitioning.clone(),
502 properties.emission_type,
503 properties.boundedness,
504 ));
505 Ok(Arc::new(Self {
506 offset: self.offset,
507 start: self.start,
508 end: self.end,
509 interval: self.interval,
510 range: self.range,
511 time_index_column: self.time_index_column.clone(),
512 time_range_column: self.time_range_column.clone(),
513 field_columns: self.field_columns.clone(),
514 output_schema: self.output_schema.clone(),
515 input: children[0].clone(),
516 metric: self.metric.clone(),
517 properties,
518 }))
519 }
520
521 fn execute(
522 &self,
523 partition: usize,
524 context: Arc<TaskContext>,
525 ) -> DataFusionResult<SendableRecordBatchStream> {
526 let baseline_metric = BaselineMetrics::new(&self.metric, partition);
527 let metrics_builder = MetricBuilder::new(&self.metric);
528 let num_series = Count::new();
529 metrics_builder
530 .with_partition(partition)
531 .build(MetricValue::Count {
532 name: METRIC_NUM_SERIES.into(),
533 count: num_series.clone(),
534 });
535
536 let input = self.input.execute(partition, context)?;
537 let schema = input.schema();
538 let time_index = schema
539 .column_with_name(&self.time_index_column)
540 .unwrap_or_else(|| panic!("time index column {} not found", self.time_index_column))
541 .0;
542 let field_columns = self
543 .field_columns
544 .iter()
545 .map(|value_col| {
546 schema
547 .column_with_name(value_col)
548 .unwrap_or_else(|| panic!("value column {value_col} not found",))
549 .0
550 })
551 .collect();
552 let time_unit = timestamp_unit(schema.field(time_index).data_type())?;
553 let aligned_ts_array =
554 RangeManipulateStream::build_aligned_ts_array(self.start, self.end, self.interval);
555 Ok(Box::pin(RangeManipulateStream {
556 offset: self.offset,
557 start: self.start,
558 end: self.end,
559 interval: self.interval,
560 range: self.range,
561 time_index,
562 time_unit,
563 field_columns,
564 aligned_ts_array,
565 output_schema: self.output_schema.clone(),
566 input,
567 metric: baseline_metric,
568 num_series,
569 }))
570 }
571
572 fn metrics(&self) -> Option<MetricsSet> {
573 Some(self.metric.clone_inner())
574 }
575
576 fn child_stats_requests(&self, partition: Option<usize>) -> Vec<ChildStats> {
577 vec![ChildStats::At(partition)]
578 }
579
580 fn statistics_from_inputs(
581 &self,
582 input_stats: &[Arc<Statistics>],
583 _args: &StatisticsArgs,
584 ) -> DataFusionResult<Arc<Statistics>> {
585 let input_stats = &input_stats[0];
586
587 let estimated_row_num = (self.end - self.start) as f64 / self.interval as f64;
588 let estimated_total_bytes = input_stats
589 .total_byte_size
590 .get_value()
591 .zip(input_stats.num_rows.get_value())
592 .map(|(size, rows)| {
593 Precision::Inexact(((*size as f64 / *rows as f64) * estimated_row_num).floor() as _)
594 })
595 .unwrap_or_default();
596
597 Ok(Arc::new(Statistics {
598 num_rows: Precision::Inexact(estimated_row_num as _),
599 total_byte_size: estimated_total_bytes,
600 column_statistics: Statistics::unknown_column(&self.schema()),
602 }))
603 }
604
605 fn name(&self) -> &str {
606 "RangeManipulateExec"
607 }
608}
609
610impl DisplayAs for RangeManipulateExec {
611 fn fmt_as(&self, t: DisplayFormatType, f: &mut std::fmt::Formatter) -> std::fmt::Result {
612 match t {
613 DisplayFormatType::Default
614 | DisplayFormatType::Verbose
615 | DisplayFormatType::TreeRender => {
616 write!(
617 f,
618 "PromRangeManipulateExec: req range=[{}..{}], interval=[{}], eval range=[{}], time index=[{}]",
619 self.start, self.end, self.interval, self.range, self.time_index_column
620 )
621 }
622 }
623 }
624}
625
626pub struct RangeManipulateStream {
627 offset: Millisecond,
628 start: Millisecond,
629 end: Millisecond,
630 interval: Millisecond,
631 range: Millisecond,
632 time_index: usize,
633 time_unit: TimeUnit,
634 field_columns: Vec<usize>,
635 aligned_ts_array: ArrayRef,
636
637 output_schema: SchemaRef,
638 input: SendableRecordBatchStream,
639 metric: BaselineMetrics,
640 num_series: Count,
642}
643
644impl RecordBatchStream for RangeManipulateStream {
645 fn schema(&self) -> SchemaRef {
646 self.output_schema.clone()
647 }
648}
649
650impl Stream for RangeManipulateStream {
651 type Item = DataFusionResult<RecordBatch>;
652
653 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
654 let poll = loop {
655 match ready!(self.input.poll_next_unpin(cx)) {
656 Some(Ok(batch)) => {
657 let timer = std::time::Instant::now();
658 let result = self.manipulate(batch);
659 if let Ok(None) = result {
660 self.metric.elapsed_compute().add_elapsed(timer);
661 continue;
662 } else {
663 self.num_series.add(1);
664 self.metric.elapsed_compute().add_elapsed(timer);
665 break Poll::Ready(result.transpose());
666 }
667 }
668 None => {
669 PROMQL_SERIES_COUNT.observe(self.num_series.value() as f64);
670 break Poll::Ready(None);
671 }
672 Some(Err(e)) => break Poll::Ready(Some(Err(e))),
673 }
674 };
675 self.metric.record_poll(poll)
676 }
677}
678
679impl RangeManipulateStream {
680 pub fn manipulate(&self, input: RecordBatch) -> DataFusionResult<Option<RecordBatch>> {
684 let mut other_columns = (0..input.columns().len()).collect::<HashSet<_>>();
685 let (ranges, (start, end)) = self.calculate_range(&input)?;
687 if ranges.iter().all(|(_, len)| *len == 0) {
689 return Ok(None);
690 }
691
692 let mut new_columns = input.columns().to_vec();
694 for index in self.field_columns.iter() {
695 let _ = other_columns.remove(index);
696 let column = input.column(*index);
697 let new_column = Arc::new(
698 RangeArray::from_ranges(column.clone(), ranges.clone())
699 .map_err(|e| ArrowError::InvalidArgumentError(e.to_string()))?
700 .into_dict(),
701 );
702 new_columns[*index] = new_column;
703 }
704
705 let scale = nanoseconds_per_native_tick(self.time_unit);
708 let (timestamps, _) = timestamp_array_to_primitive(input.column(self.time_index))
709 .ok_or_else(|| {
710 DataFusionError::Execution("Time index column is not a timestamp".into())
711 })?;
712 let timestamp_values = timestamps
713 .values()
714 .iter()
715 .enumerate()
716 .map(|(index, timestamp)| {
717 if !input.column(self.time_index).is_valid(index) {
718 return Ok(None);
719 }
720 let shifted_ns = (*timestamp as i128) * scale + (self.offset as i128) * 1_000_000;
721 i64::try_from(shifted_ns / 1_000_000)
722 .map(Some)
723 .map_err(|_| {
724 ArrowError::ComputeError(
725 "RangeManipulate timestamp payload overflow".into(),
726 )
727 })
728 })
729 .collect::<std::result::Result<Vec<_>, _>>()?;
730 let timestamp_values = TimestampMillisecondArray::from(timestamp_values);
731 let ts_range_column = RangeArray::from_ranges(Arc::new(timestamp_values), ranges.clone())
732 .map_err(|e| ArrowError::InvalidArgumentError(e.to_string()))?
733 .into_dict();
734 new_columns.push(Arc::new(ts_range_column));
735
736 let take_indices = Int64Array::from(vec![0; ranges.len()]);
738 for index in other_columns.into_iter() {
739 new_columns[index] = compute::take(&input.column(index), &take_indices, None)?;
740 }
741 let new_time_index = if ranges.len() != self.aligned_ts_array.len() {
743 Self::build_aligned_ts_array(start, end, self.interval)
744 } else {
745 self.aligned_ts_array.clone()
746 };
747 new_columns[self.time_index] = new_time_index;
748
749 RecordBatch::try_new(self.output_schema.clone(), new_columns)
750 .map(Some)
751 .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))
752 }
753
754 fn build_aligned_ts_array(start: i64, end: i64, interval: i64) -> ArrayRef {
755 Arc::new(TimestampMillisecondArray::from_iter_values(
756 (start..=end).step_by(interval as _),
757 ))
758 }
759
760 #[allow(clippy::type_complexity)]
764 fn calculate_range(
765 &self,
766 input: &RecordBatch,
767 ) -> DataFusionResult<(Vec<(u32, u32)>, (i64, i64))> {
768 let ts_column = input.column(self.time_index);
769 let scale = nanoseconds_per_native_tick(self.time_unit);
770 let (timestamps, _) = timestamp_array_to_primitive(ts_column).ok_or_else(|| {
771 DataFusionError::Execution("Time index column is not a timestamp".into())
772 })?;
773 let timestamps = timestamps.values();
774 let timestamp =
775 |index| (timestamps[index] as i128) * scale + (self.offset as i128) * 1_000_000;
776 let len = timestamps.len();
777 if len == 0 {
778 return Ok((vec![], (self.start, self.end)));
779 }
780
781 let query_start = self.start as i128;
784 let query_end = self.end as i128;
785 let interval = self.interval as i128;
786 let first_ts = timestamp(0).div_euclid(1_000_000);
787 let remainder = (first_ts - query_start).rem_euclid(interval);
789 let first_ts_aligned = first_ts + (interval - remainder).rem_euclid(interval);
790 let last_ts_with_range =
791 (timestamp(len - 1) + (self.range as i128) * 1_000_000).div_euclid(1_000_000);
792 let remainder = (last_ts_with_range - query_start).rem_euclid(interval);
793 let last_ts_aligned = last_ts_with_range - remainder;
794 let start = query_start.max(first_ts_aligned);
795 let end = query_end.min(last_ts_aligned);
796 if start > end {
797 return Ok((vec![], (self.start, self.end)));
798 }
799 let start = start as i64;
801 let end = end as i64;
802 let mut ranges = Vec::new();
803
804 let mut left = 0usize;
811 let mut right = 0usize;
812 for curr_ts in (start..=end).step_by(self.interval as _) {
813 let start_ts = (curr_ts as i128) * 1_000_000 - (self.range as i128) * 1_000_000;
814
815 while left < len && timestamp(left) <= start_ts {
816 left += 1;
817 }
818 right = right.max(left);
819 while right < len && timestamp(right) <= (curr_ts as i128) * 1_000_000 {
820 right += 1;
821 }
822
823 if left == right {
824 ranges.push((0, 0));
825 } else {
826 ranges.push((left as _, (right - left) as _));
827 }
828 }
829
830 Ok((ranges, (start, end)))
831 }
832}
833
834#[cfg(test)]
835mod test {
836 use datafusion::arrow::array::{
837 ArrayRef, DictionaryArray, Float64Array, StringArray, TimestampMicrosecondArray,
838 TimestampNanosecondArray, TimestampSecondArray,
839 };
840 use datafusion::arrow::buffer::NullBuffer;
841 use datafusion::arrow::datatypes::{
842 ArrowPrimitiveType, DataType, Field, Int64Type, Schema, TimestampMillisecondType,
843 };
844 use datafusion::common::ToDFSchema;
845 use datafusion::datasource::memory::MemorySourceConfig;
846 use datafusion::datasource::source::DataSourceExec;
847 use datafusion::logical_expr::{
848 EmptyRelation, Extension, LogicalPlan, UserDefinedLogicalNodeCore,
849 };
850 use datafusion::physical_expr::Partitioning;
851 use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
852 use datafusion::physical_plan::memory::MemoryStream;
853 use datafusion::physical_plan::{ChildrenPropertiesMode, ReplaceChildrenOptions};
854 use datafusion::prelude::SessionContext;
855 use datatypes::arrow::array::TimestampMillisecondArray;
856 use futures::FutureExt;
857
858 use super::*;
859
860 const TIME_INDEX_COLUMN: &str = "timestamp";
861
862 fn project_batch(batch: &RecordBatch, indices: &[usize]) -> RecordBatch {
863 let fields = indices
864 .iter()
865 .map(|&idx| batch.schema().field(idx).clone())
866 .collect::<Vec<_>>();
867 let columns = indices
868 .iter()
869 .map(|&idx| batch.column(idx).clone())
870 .collect::<Vec<_>>();
871 let schema = Arc::new(Schema::new(fields));
872 RecordBatch::try_new(schema, columns).unwrap()
873 }
874
875 fn prepare_test_data() -> DataSourceExec {
876 let schema = Arc::new(Schema::new(vec![
877 Field::new(TIME_INDEX_COLUMN, TimestampMillisecondType::DATA_TYPE, true),
878 Field::new("value_1", DataType::Float64, true),
879 Field::new("value_2", DataType::Float64, true),
880 Field::new("path", DataType::Utf8, true),
881 ]));
882 let timestamp_column = Arc::new(TimestampMillisecondArray::from(vec![
883 0, 30_000, 60_000, 90_000, 120_000, 180_000, 240_000, 241_000, 271_000, 291_000, ])) as _;
887 let field_column: ArrayRef = Arc::new(Float64Array::from(vec![1.0; 10])) as _;
888 let path_column = Arc::new(StringArray::from(vec!["foo"; 10])) as _;
889 let data = RecordBatch::try_new(
890 schema.clone(),
891 vec![
892 timestamp_column,
893 field_column.clone(),
894 field_column,
895 path_column,
896 ],
897 )
898 .unwrap();
899
900 DataSourceExec::new(Arc::new(
901 MemorySourceConfig::try_new(&[vec![data]], schema, None).unwrap(),
902 ))
903 }
904
905 async fn do_normalize_test(
906 start: Millisecond,
907 end: Millisecond,
908 interval: Millisecond,
909 range: Millisecond,
910 expected: String,
911 ) {
912 let memory_exec = Arc::new(prepare_test_data());
913 let time_index = TIME_INDEX_COLUMN.to_string();
914 let field_columns = vec!["value_1".to_string(), "value_2".to_string()];
915 let manipulate_output_schema = SchemaRef::new(
916 RangeManipulate::calculate_output_schema(
917 &memory_exec.schema().to_dfschema_ref().unwrap(),
918 &time_index,
919 &field_columns,
920 )
921 .unwrap()
922 .as_arrow()
923 .clone(),
924 );
925 let properties = Arc::new(PlanProperties::new(
926 EquivalenceProperties::new(manipulate_output_schema.clone()),
927 Partitioning::UnknownPartitioning(1),
928 EmissionType::Incremental,
929 Boundedness::Bounded,
930 ));
931 let normalize_exec = Arc::new(RangeManipulateExec {
932 offset: 0,
933 start,
934 end,
935 interval,
936 range,
937 field_columns,
938 output_schema: manipulate_output_schema,
939 time_range_column: RangeManipulate::build_timestamp_range_name(&time_index),
940 time_index_column: time_index,
941 input: memory_exec,
942 metric: ExecutionPlanMetricsSet::new(),
943 properties,
944 });
945 let session_context = SessionContext::default();
946 let result = datafusion::physical_plan::collect(normalize_exec, session_context.task_ctx())
947 .await
948 .unwrap();
949 let result_literal: String = result
951 .into_iter()
952 .filter_map(|batch| {
953 batch
954 .columns()
955 .iter()
956 .map(|array| {
957 if matches!(array.data_type(), &DataType::Dictionary(..)) {
958 let dict_array = array
959 .as_any()
960 .downcast_ref::<DictionaryArray<Int64Type>>()
961 .unwrap()
962 .clone();
963 format!("{:?}", RangeArray::try_new(dict_array).unwrap())
964 } else {
965 format!("{array:?}")
966 }
967 })
968 .reduce(|lhs, rhs| lhs + "\n" + &rhs)
969 })
970 .reduce(|lhs, rhs| lhs + "\n\n" + &rhs)
971 .unwrap();
972
973 assert_eq!(result_literal, expected);
974 }
975
976 #[tokio::test]
977 async fn native_timestamps_preserve_range_membership_and_ms_payload() {
978 for (unit, ticks_per_ms) in [
979 (TimeUnit::Microsecond, 1_000_i64),
980 (TimeUnit::Nanosecond, 1_000_000_i64),
981 ] {
982 let lower = 1_000 * ticks_per_ms;
983 let upper = 1_001 * ticks_per_ms;
984 let timestamps = vec![lower, lower + 1, lower + 2, upper, upper + 1];
987 let time: ArrayRef = match unit {
988 TimeUnit::Microsecond => Arc::new(TimestampMicrosecondArray::from(timestamps)),
989 TimeUnit::Nanosecond => Arc::new(TimestampNanosecondArray::from(timestamps)),
990 _ => unreachable!(),
991 };
992 let schema = Arc::new(Schema::new(vec![
993 Field::new(TIME_INDEX_COLUMN, DataType::Timestamp(unit, None), false),
994 Field::new("value", DataType::Float64, true),
995 ]));
996 let batch = RecordBatch::try_new(
997 schema.clone(),
998 vec![
999 time,
1000 Arc::new(Float64Array::from(vec![10.0, 20.0, 30.0, 40.0, 50.0])),
1001 ],
1002 )
1003 .unwrap();
1004 let logical_input = LogicalPlan::EmptyRelation(EmptyRelation {
1005 produce_one_row: false,
1006 schema: schema.clone().to_dfschema_ref().unwrap(),
1007 });
1008 let plan = RangeManipulate::new(
1009 1_001,
1010 1_001,
1011 1,
1012 0,
1013 1,
1014 TIME_INDEX_COLUMN.to_string(),
1015 vec!["value".to_string()],
1016 logical_input.clone(),
1017 )
1018 .unwrap();
1019 let output_time = Field::new(
1020 TIME_INDEX_COLUMN,
1021 DataType::Timestamp(TimeUnit::Millisecond, None),
1022 false,
1023 );
1024 let output_schema = Arc::new(Schema::new(vec![
1025 output_time.clone(),
1026 RangeArray::convert_field(&Field::new("value", DataType::Float64, true)),
1027 Field::new(
1028 RangeManipulate::build_timestamp_range_name(TIME_INDEX_COLUMN),
1029 RangeArray::convert_field(&output_time).data_type().clone(),
1030 false,
1031 ),
1032 ]));
1033 assert_eq!(plan.schema().as_arrow(), output_schema.as_ref());
1034
1035 let rebuilt = RangeManipulate::deserialize(&plan.serialize())
1036 .unwrap()
1037 .with_exprs_and_inputs(vec![], vec![logical_input])
1038 .unwrap();
1039 assert_eq!(rebuilt.schema(), plan.schema());
1040 assert_eq!(rebuilt.input.schema().as_arrow(), schema.as_ref());
1041
1042 let input = Arc::new(DataSourceExec::new(Arc::new(
1043 MemorySourceConfig::try_new(&[vec![batch]], schema.clone(), None).unwrap(),
1044 )));
1045 let exec = rebuilt.to_execution_plan(input);
1046 assert_eq!(exec.schema(), output_schema);
1047 assert_eq!(exec.children()[0].schema(), schema);
1048
1049 let batches =
1050 datafusion::physical_plan::collect(exec, SessionContext::default().task_ctx())
1051 .await
1052 .unwrap();
1053 assert_eq!(batches.len(), 1, "{unit:?}");
1054 let output = &batches[0];
1055 assert_eq!(output.schema(), output_schema);
1056 assert_eq!(output.num_rows(), 1);
1057 assert_eq!(
1058 output
1059 .column(0)
1060 .as_any()
1061 .downcast_ref::<TimestampMillisecondArray>()
1062 .unwrap()
1063 .values()
1064 .as_ref(),
1065 &[1_001]
1066 );
1067
1068 let values = RangeArray::try_new(
1071 output
1072 .column(1)
1073 .as_any()
1074 .downcast_ref::<DictionaryArray<Int64Type>>()
1075 .unwrap()
1076 .clone(),
1077 )
1078 .unwrap();
1079 assert_eq!(values.get_offset_length(0), Some((1, 3)));
1080 assert_eq!(
1081 values.get(0).unwrap().to_data(),
1082 Float64Array::from(vec![20.0, 30.0, 40.0]).to_data()
1083 );
1084 let timestamps = RangeArray::try_new(
1085 output
1086 .column(2)
1087 .as_any()
1088 .downcast_ref::<DictionaryArray<Int64Type>>()
1089 .unwrap()
1090 .clone(),
1091 )
1092 .unwrap();
1093 assert_eq!(timestamps.get_offset_length(0), Some((1, 3)));
1094 assert_eq!(
1095 timestamps.get(0).unwrap().to_data(),
1096 TimestampMillisecondArray::from(vec![1_000, 1_000, 1_001]).to_data()
1097 );
1098 }
1099 }
1100
1101 #[test]
1102 fn logical_offset_participates_in_ordering() {
1103 let input = LogicalPlan::EmptyRelation(EmptyRelation {
1104 produce_one_row: false,
1105 schema: prepare_test_data().schema().to_dfschema_ref().unwrap(),
1106 });
1107 let first = RangeManipulate::new(
1108 0,
1109 0,
1110 0,
1111 0,
1112 0,
1113 TIME_INDEX_COLUMN.to_string(),
1114 vec!["value_1".to_string()],
1115 input.clone(),
1116 )
1117 .unwrap();
1118 let second = RangeManipulate::new(
1119 0,
1120 0,
1121 0,
1122 1,
1123 0,
1124 TIME_INDEX_COLUMN.to_string(),
1125 vec!["value_1".to_string()],
1126 input,
1127 )
1128 .unwrap();
1129 assert_ne!(first, second);
1130 assert_eq!(first.partial_cmp(&second), Some(std::cmp::Ordering::Less));
1131 }
1132
1133 #[tokio::test]
1134 async fn logical_normalize_offset_survives_rebuild_and_executes() {
1135 for (name, time_unit, raw, offset, start, range, expected_payload) in [
1136 (
1137 "millisecond offset",
1138 TimeUnit::Millisecond,
1139 0,
1140 1_000,
1141 1_000,
1142 1_000,
1143 1_000,
1144 ),
1145 (
1146 "negative native lower limit with positive window",
1147 TimeUnit::Nanosecond,
1148 -9_223_112_837_000_000_000,
1149 -259_200_000,
1150 -9_223_372_037_000,
1151 300_000,
1152 -9_223_372_037_000,
1153 ),
1154 (
1155 "second timestamp with negative fractional offset",
1156 TimeUnit::Second,
1157 1,
1158 -500,
1159 1_000,
1160 1_000,
1161 500,
1162 ),
1163 ] {
1164 let schema = Arc::new(Schema::new(vec![
1165 Field::new(
1166 TIME_INDEX_COLUMN,
1167 DataType::Timestamp(time_unit, None),
1168 false,
1169 ),
1170 Field::new("value", DataType::Float64, true),
1171 ]));
1172 let input = LogicalPlan::EmptyRelation(EmptyRelation {
1173 produce_one_row: false,
1174 schema: schema.clone().to_dfschema_ref().unwrap(),
1175 });
1176 let normalize = crate::extension_plan::SeriesNormalize::new(
1177 offset,
1178 TIME_INDEX_COLUMN,
1179 false,
1180 Vec::new(),
1181 input.clone(),
1182 );
1183 let normalize =
1184 crate::extension_plan::SeriesNormalize::deserialize(&normalize.serialize())
1185 .unwrap()
1186 .with_exprs_and_inputs(vec![], vec![input.clone()])
1187 .unwrap();
1188 let normalized = LogicalPlan::Extension(Extension {
1189 node: Arc::new(normalize),
1190 });
1191 let fresh = RangeManipulate::new(
1192 start,
1193 start,
1194 1,
1195 offset,
1196 range,
1197 TIME_INDEX_COLUMN.to_string(),
1198 vec!["value".to_string()],
1199 input.clone(),
1200 )
1201 .unwrap()
1202 .with_exprs_and_inputs(vec![], vec![input.clone()])
1203 .unwrap();
1204 let serialized = RangeManipulate::new(
1205 start,
1206 start,
1207 1,
1208 offset,
1209 range,
1210 TIME_INDEX_COLUMN.to_string(),
1211 vec!["value".to_string()],
1212 normalized.clone(),
1213 )
1214 .unwrap();
1215 let decoded = RangeManipulate::deserialize(&serialized.serialize())
1216 .unwrap()
1217 .with_exprs_and_inputs(vec![], vec![normalized])
1218 .unwrap();
1219 let timestamp: ArrayRef = match time_unit {
1220 TimeUnit::Millisecond => Arc::new(TimestampMillisecondArray::from(vec![raw])),
1221 TimeUnit::Nanosecond => Arc::new(TimestampNanosecondArray::from(vec![raw])),
1222 TimeUnit::Second => Arc::new(TimestampSecondArray::from(vec![raw])),
1223 _ => unreachable!(),
1224 };
1225 let batch = RecordBatch::try_new(
1226 schema.clone(),
1227 vec![timestamp, Arc::new(Float64Array::from(vec![7.0]))],
1228 )
1229 .unwrap();
1230 for (mode, rebuilt) in [("fresh", fresh), ("decoded", decoded)] {
1231 let rebuilt = rebuilt
1232 .with_exprs_and_inputs(vec![], vec![input.clone()])
1233 .unwrap();
1234 assert_eq!(rebuilt.offset, offset, "{name}: {mode}");
1235 assert_eq!(rebuilt.input.schema(), input.schema(), "{name}: {mode}");
1236
1237 let empty_exec_input = Arc::new(DataSourceExec::new(Arc::new(
1238 MemorySourceConfig::try_new(&[vec![]], schema.clone(), None).unwrap(),
1239 )));
1240 let exec_input = Arc::new(DataSourceExec::new(Arc::new(
1241 MemorySourceConfig::try_new(&[vec![batch.clone()]], schema.clone(), None)
1242 .unwrap(),
1243 )));
1244 let exec = rebuilt
1245 .to_execution_plan(empty_exec_input)
1246 .replace_children(
1247 vec![exec_input],
1248 ReplaceChildrenOptions::new(ChildrenPropertiesMode::Recompute),
1249 )
1250 .unwrap();
1251 let output =
1252 datafusion::physical_plan::collect(exec, SessionContext::default().task_ctx())
1253 .await
1254 .unwrap();
1255 assert_eq!(output.len(), 1, "{name}: {mode}");
1256 let output = &output[0];
1257 assert_eq!(output.num_rows(), 1, "{name}: {mode}");
1258 assert_eq!(
1259 output
1260 .column(0)
1261 .as_any()
1262 .downcast_ref::<TimestampMillisecondArray>()
1263 .unwrap()
1264 .value(0),
1265 start,
1266 "{name}: {mode}"
1267 );
1268 let values = RangeArray::try_new(
1269 output
1270 .column(1)
1271 .as_any()
1272 .downcast_ref::<DictionaryArray<Int64Type>>()
1273 .unwrap()
1274 .clone(),
1275 )
1276 .unwrap();
1277 assert_eq!(values.get_offset_length(0), Some((0, 1)), "{name}: {mode}");
1278 assert_eq!(
1279 values.get(0).unwrap().to_data(),
1280 Float64Array::from(vec![7.0]).to_data(),
1281 "{name}: {mode}"
1282 );
1283 let timestamps = RangeArray::try_new(
1284 output
1285 .column(2)
1286 .as_any()
1287 .downcast_ref::<DictionaryArray<Int64Type>>()
1288 .unwrap()
1289 .clone(),
1290 )
1291 .unwrap();
1292 assert_eq!(
1293 timestamps.get_offset_length(0),
1294 Some((0, 1)),
1295 "{name}: {mode}"
1296 );
1297 assert_eq!(
1298 timestamps.get(0).unwrap().to_data(),
1299 TimestampMillisecondArray::from(vec![expected_payload]).to_data(),
1300 "{name}: {mode}"
1301 );
1302 }
1303 }
1304 }
1305
1306 #[tokio::test]
1307 async fn range_payload_preserves_null_timestamp_and_rejects_offset_overflow() {
1308 let schema = Arc::new(Schema::new(vec![
1309 Field::new(TIME_INDEX_COLUMN, TimestampMillisecondType::DATA_TYPE, true),
1310 Field::new("value", DataType::Float64, true),
1311 ]));
1312 let null_timestamp =
1313 TimestampMillisecondArray::new(vec![1_000].into(), Some(NullBuffer::from(vec![false])));
1314 let batch = RecordBatch::try_new(
1315 schema.clone(),
1316 vec![
1317 Arc::new(null_timestamp),
1318 Arc::new(Float64Array::from(vec![7.0])),
1319 ],
1320 )
1321 .unwrap();
1322 let input = Arc::new(DataSourceExec::new(Arc::new(
1323 MemorySourceConfig::try_new(&[vec![batch]], schema.clone(), None).unwrap(),
1324 )));
1325 let plan = RangeManipulate::new(
1326 1_000,
1327 1_000,
1328 1,
1329 0,
1330 1,
1331 TIME_INDEX_COLUMN.to_string(),
1332 vec!["value".to_string()],
1333 LogicalPlan::EmptyRelation(EmptyRelation {
1334 produce_one_row: false,
1335 schema: schema.clone().to_dfschema_ref().unwrap(),
1336 }),
1337 )
1338 .unwrap();
1339 let output = datafusion::physical_plan::collect(
1340 plan.to_execution_plan(input),
1341 SessionContext::default().task_ctx(),
1342 )
1343 .await
1344 .unwrap();
1345 let timestamps = RangeArray::try_new(
1346 output[0]
1347 .column(2)
1348 .as_any()
1349 .downcast_ref::<DictionaryArray<Int64Type>>()
1350 .unwrap()
1351 .clone(),
1352 )
1353 .unwrap();
1354 let payload = timestamps.get(0).unwrap();
1355 let payload = payload
1356 .as_any()
1357 .downcast_ref::<TimestampMillisecondArray>()
1358 .unwrap();
1359 assert_eq!(payload.len(), 1);
1360 assert!(!payload.is_valid(0));
1361
1362 let batch = RecordBatch::try_new(
1363 schema.clone(),
1364 vec![
1365 Arc::new(TimestampMillisecondArray::from(vec![0, i64::MAX])),
1366 Arc::new(Float64Array::from(vec![7.0, 8.0])),
1367 ],
1368 )
1369 .unwrap();
1370 let input = Arc::new(DataSourceExec::new(Arc::new(
1371 MemorySourceConfig::try_new(&[vec![batch]], schema.clone(), None).unwrap(),
1372 )));
1373 let normalized = crate::extension_plan::SeriesNormalize::new(
1374 1,
1375 TIME_INDEX_COLUMN,
1376 false,
1377 Vec::new(),
1378 LogicalPlan::EmptyRelation(EmptyRelation {
1379 produce_one_row: false,
1380 schema: schema.to_dfschema_ref().unwrap(),
1381 }),
1382 );
1383 let plan = RangeManipulate::new(
1384 1,
1385 1,
1386 1,
1387 1,
1388 1,
1389 TIME_INDEX_COLUMN.to_string(),
1390 vec!["value".to_string()],
1391 LogicalPlan::Extension(Extension {
1392 node: Arc::new(normalized),
1393 }),
1394 )
1395 .unwrap();
1396 let error = datafusion::physical_plan::collect(
1397 plan.to_execution_plan(input),
1398 SessionContext::default().task_ctx(),
1399 )
1400 .await
1401 .unwrap_err();
1402 assert!(error.to_string().contains("timestamp payload overflow"));
1403 }
1404
1405 #[tokio::test]
1406 async fn pruning_should_keep_time_and_value_columns_for_exec() {
1407 let schema = Arc::new(Schema::new(vec![
1408 Field::new(TIME_INDEX_COLUMN, TimestampMillisecondType::DATA_TYPE, true),
1409 Field::new("value_1", DataType::Float64, true),
1410 Field::new("value_2", DataType::Float64, true),
1411 Field::new("path", DataType::Utf8, true),
1412 ]));
1413 let df_schema = schema.clone().to_dfschema_ref().unwrap();
1414 let input = LogicalPlan::EmptyRelation(EmptyRelation {
1415 produce_one_row: false,
1416 schema: df_schema,
1417 });
1418 let plan = RangeManipulate::new(
1419 0,
1420 310_000,
1421 30_000,
1422 0,
1423 90_000,
1424 TIME_INDEX_COLUMN.to_string(),
1425 vec!["value_1".to_string(), "value_2".to_string()],
1426 input,
1427 )
1428 .unwrap();
1429
1430 let output_columns = [3usize];
1432 let required = plan.necessary_children_exprs(&output_columns).unwrap();
1433 let required = &required[0];
1434 assert_eq!(required.as_slice(), &[0, 1, 2, 3]);
1435
1436 let timestamp_column = Arc::new(TimestampMillisecondArray::from(vec![
1437 0, 30_000, 60_000, 90_000, 120_000, 180_000, 240_000, 241_000, 271_000, 291_000, ])) as _;
1441 let field_column: ArrayRef = Arc::new(Float64Array::from(vec![1.0; 10])) as _;
1442 let path_column = Arc::new(StringArray::from(vec!["foo"; 10])) as _;
1443 let input_batch = RecordBatch::try_new(
1444 schema,
1445 vec![
1446 timestamp_column,
1447 field_column.clone(),
1448 field_column,
1449 path_column,
1450 ],
1451 )
1452 .unwrap();
1453
1454 let projected = project_batch(&input_batch, required);
1455 let projected_schema = projected.schema();
1456 let memory_exec = Arc::new(DataSourceExec::new(Arc::new(
1457 MemorySourceConfig::try_new(&[vec![projected]], projected_schema, None).unwrap(),
1458 )));
1459 let range_exec = plan.to_execution_plan(memory_exec);
1460 let session_context = SessionContext::default();
1461 let output_batches =
1462 datafusion::physical_plan::collect(range_exec, session_context.task_ctx())
1463 .await
1464 .unwrap();
1465 assert_eq!(output_batches.len(), 1);
1466
1467 let output_batch = &output_batches[0];
1468 let path = output_batch
1469 .column(3)
1470 .as_any()
1471 .downcast_ref::<StringArray>()
1472 .unwrap();
1473 assert!(path.iter().all(|v| v == Some("foo")));
1474
1475 let broken_required = [3usize];
1477 let broken = project_batch(&input_batch, &broken_required);
1478 let broken_schema = broken.schema();
1479 let broken_exec = Arc::new(DataSourceExec::new(Arc::new(
1480 MemorySourceConfig::try_new(&[vec![broken]], broken_schema, None).unwrap(),
1481 )));
1482 let broken_range_exec = plan.to_execution_plan(broken_exec);
1483 let session_context = SessionContext::default();
1484 let broken_result = std::panic::AssertUnwindSafe(async {
1485 datafusion::physical_plan::collect(broken_range_exec, session_context.task_ctx()).await
1486 })
1487 .catch_unwind()
1488 .await;
1489 assert!(broken_result.is_err());
1490 }
1491
1492 #[tokio::test]
1493 async fn interval_30s_range_90s() {
1494 let expected = String::from(
1495 "PrimitiveArray<Timestamp(ms)>\n[\n \
1496 1970-01-01T00:00:00,\n \
1497 1970-01-01T00:00:30,\n \
1498 1970-01-01T00:01:00,\n \
1499 1970-01-01T00:01:30,\n \
1500 1970-01-01T00:02:00,\n \
1501 1970-01-01T00:02:30,\n \
1502 1970-01-01T00:03:00,\n \
1503 1970-01-01T00:03:30,\n \
1504 1970-01-01T00:04:00,\n \
1505 1970-01-01T00:04:30,\n \
1506 1970-01-01T00:05:00,\n\
1507 ]\nRangeArray { \
1508 base array: PrimitiveArray<Float64>\n[\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n], \
1509 ranges: [Some(0..1), Some(0..2), Some(0..3), Some(1..4), Some(2..5), Some(3..5), Some(4..6), Some(5..6), Some(5..7), Some(6..8), Some(6..10)] \
1510 }\nRangeArray { \
1511 base array: PrimitiveArray<Float64>\n[\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n], \
1512 ranges: [Some(0..1), Some(0..2), Some(0..3), Some(1..4), Some(2..5), Some(3..5), Some(4..6), Some(5..6), Some(5..7), Some(6..8), Some(6..10)] \
1513 }\nStringArray\n[\n \"foo\",\n \"foo\",\n \"foo\",\n \"foo\",\n \"foo\",\n \"foo\",\n \"foo\",\n \"foo\",\n \"foo\",\n \"foo\",\n \"foo\",\n]\n\
1514 RangeArray { \
1515 base array: PrimitiveArray<Timestamp(ms)>\n[\n 1970-01-01T00:00:00,\n 1970-01-01T00:00:30,\n 1970-01-01T00:01:00,\n 1970-01-01T00:01:30,\n 1970-01-01T00:02:00,\n 1970-01-01T00:03:00,\n 1970-01-01T00:04:00,\n 1970-01-01T00:04:01,\n 1970-01-01T00:04:31,\n 1970-01-01T00:04:51,\n], \
1516 ranges: [Some(0..1), Some(0..2), Some(0..3), Some(1..4), Some(2..5), Some(3..5), Some(4..6), Some(5..6), Some(5..7), Some(6..8), Some(6..10)] \
1517 }",
1518 );
1519 do_normalize_test(0, 310_000, 30_000, 90_000, expected.clone()).await;
1520
1521 do_normalize_test(-300000, 310_000, 30_000, 90_000, expected).await;
1523 }
1524
1525 #[tokio::test]
1526 async fn small_empty_range() {
1527 let expected = String::from(
1528 "PrimitiveArray<Timestamp(ms)>\n[\n \
1529 1970-01-01T00:00:00.001,\n \
1530 1970-01-01T00:00:03.001,\n \
1531 1970-01-01T00:00:06.001,\n \
1532 1970-01-01T00:00:09.001,\n\
1533 ]\nRangeArray { \
1534 base array: PrimitiveArray<Float64>\n[\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n], \
1535 ranges: [Some(0..1), Some(0..0), Some(0..0), Some(0..0)] \
1536 }\nRangeArray { \
1537 base array: PrimitiveArray<Float64>\n[\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n 1.0,\n], \
1538 ranges: [Some(0..1), Some(0..0), Some(0..0), Some(0..0)] \
1539 }\nStringArray\n[\n \"foo\",\n \"foo\",\n \"foo\",\n \"foo\",\n]\n\
1540 RangeArray { \
1541 base array: PrimitiveArray<Timestamp(ms)>\n[\n 1970-01-01T00:00:00,\n 1970-01-01T00:00:30,\n 1970-01-01T00:01:00,\n 1970-01-01T00:01:30,\n 1970-01-01T00:02:00,\n 1970-01-01T00:03:00,\n 1970-01-01T00:04:00,\n 1970-01-01T00:04:01,\n 1970-01-01T00:04:31,\n 1970-01-01T00:04:51,\n], \
1542 ranges: [Some(0..1), Some(0..0), Some(0..0), Some(0..0)] \
1543 }",
1544 );
1545 do_normalize_test(1, 10_001, 3_000, 1_000, expected).await;
1546 }
1547
1548 #[test]
1549 fn test_calculate_range_preserves_alignment() {
1550 let schema = Arc::new(Schema::new(vec![Field::new(
1553 "timestamp",
1554 TimestampMillisecondType::DATA_TYPE,
1555 false,
1556 )]));
1557 let empty_stream = MemoryStream::try_new(vec![], schema.clone(), None).unwrap();
1558
1559 let stream = RangeManipulateStream {
1560 offset: 0,
1561 start: 1758093274000, end: 1758093334000, interval: 30000, range: 60000, time_index: 0,
1566 time_unit: TimeUnit::Millisecond,
1567 field_columns: vec![],
1568 aligned_ts_array: Arc::new(TimestampMillisecondArray::from(vec![0i64; 0])),
1569 output_schema: schema.clone(),
1570 input: Box::pin(empty_stream),
1571 metric: BaselineMetrics::new(&ExecutionPlanMetricsSet::new(), 0),
1572 num_series: Count::new(),
1573 };
1574
1575 let test_timestamps = vec![
1577 1758093260000, 1758093290000, 1758093320000, ];
1581 let ts_array = TimestampMillisecondArray::from(test_timestamps);
1582 let test_schema = Arc::new(Schema::new(vec![Field::new(
1583 "timestamp",
1584 TimestampMillisecondType::DATA_TYPE,
1585 false,
1586 )]));
1587 let batch = RecordBatch::try_new(test_schema, vec![Arc::new(ts_array)]).unwrap();
1588
1589 let (ranges, (start, end)) = stream.calculate_range(&batch).unwrap();
1590
1591 assert_eq!(
1593 start % 30000,
1594 1758093274000 % 30000,
1595 "Optimized start should preserve query alignment pattern"
1596 );
1597
1598 let expected_timestamps: Vec<i64> = (start..=end).step_by(30000).collect();
1600 assert_eq!(ranges.len(), expected_timestamps.len());
1601
1602 for ts in expected_timestamps {
1604 assert_eq!(
1605 ts % 30000,
1606 1758093274000 % 30000,
1607 "All timestamps should maintain query alignment pattern"
1608 );
1609 }
1610 }
1611
1612 #[tokio::test]
1613 async fn no_intersection_batch_is_skipped_and_stream_continues() {
1614 let schema = Arc::new(Schema::new(vec![
1615 Field::new(
1616 TIME_INDEX_COLUMN,
1617 TimestampMillisecondType::DATA_TYPE,
1618 false,
1619 ),
1620 Field::new("value", DataType::Float64, false),
1621 ]));
1622 let input = LogicalPlan::EmptyRelation(EmptyRelation {
1623 produce_one_row: false,
1624 schema: schema.clone().to_dfschema_ref().unwrap(),
1625 });
1626 let plan = RangeManipulate::new(
1627 0,
1628 50,
1629 10,
1630 0,
1631 1,
1632 TIME_INDEX_COLUMN.to_string(),
1633 vec!["value".to_string()],
1634 input,
1635 )
1636 .unwrap();
1637 let no_intersection = RecordBatch::try_new(
1638 schema.clone(),
1639 vec![
1640 Arc::new(TimestampMillisecondArray::from(vec![100])),
1641 Arc::new(Float64Array::from(vec![1.0])),
1642 ],
1643 )
1644 .unwrap();
1645 let intersection = RecordBatch::try_new(
1646 schema.clone(),
1647 vec![
1648 Arc::new(TimestampMillisecondArray::from(vec![20])),
1649 Arc::new(Float64Array::from(vec![2.0])),
1650 ],
1651 )
1652 .unwrap();
1653 let input = Arc::new(DataSourceExec::new(Arc::new(
1654 MemorySourceConfig::try_new(&[vec![no_intersection, intersection]], schema, None)
1655 .unwrap(),
1656 )));
1657
1658 let batches = datafusion::physical_plan::collect(
1659 plan.to_execution_plan(input),
1660 SessionContext::default().task_ctx(),
1661 )
1662 .await
1663 .unwrap();
1664
1665 assert_eq!(batches.len(), 1);
1666 assert_eq!(batches[0].num_rows(), 1);
1667 let timestamps = batches[0]
1668 .column(0)
1669 .as_any()
1670 .downcast_ref::<TimestampMillisecondArray>()
1671 .unwrap();
1672 assert_eq!(timestamps.value(0), 20);
1673 }
1674
1675 fn calculate_range_for_test(
1676 query_start: i64,
1677 query_end: i64,
1678 interval: i64,
1679 range: i64,
1680 timestamps: &[i64],
1681 ) -> (Vec<(u32, u32)>, (i64, i64)) {
1682 let schema = Arc::new(Schema::new(vec![Field::new(
1683 TIME_INDEX_COLUMN,
1684 TimestampMillisecondType::DATA_TYPE,
1685 false,
1686 )]));
1687 let empty_stream = MemoryStream::try_new(vec![], schema.clone(), None).unwrap();
1688 let stream = RangeManipulateStream {
1689 offset: 0,
1690 start: query_start,
1691 end: query_end,
1692 interval,
1693 range,
1694 time_index: 0,
1695 time_unit: TimeUnit::Millisecond,
1696 field_columns: vec![],
1697 aligned_ts_array: Arc::new(TimestampMillisecondArray::from(vec![0i64; 0])),
1698 output_schema: schema.clone(),
1699 input: Box::pin(empty_stream),
1700 metric: BaselineMetrics::new(&ExecutionPlanMetricsSet::new(), 0),
1701 num_series: Count::new(),
1702 };
1703 let batch = RecordBatch::try_new(
1704 schema,
1705 vec![Arc::new(TimestampMillisecondArray::from(
1706 timestamps.to_vec(),
1707 ))],
1708 )
1709 .unwrap();
1710
1711 stream.calculate_range(&batch).unwrap()
1712 }
1713
1714 #[test]
1715 fn calculate_range_keeps_query_aligned_tail() {
1716 let (ranges, bounds) = calculate_range_for_test(4, 94, 30, 15, &[20, 50, 80]);
1717
1718 assert_eq!(bounds, (34, 94));
1719 assert_eq!(ranges, vec![(0, 1), (1, 1), (2, 1)]);
1720 }
1721
1722 fn calculate_range_oracle(
1727 timestamps: &[i64],
1728 start: i64,
1729 end: i64,
1730 interval: i64,
1731 range: i64,
1732 ) -> Vec<(u32, u32)> {
1733 if timestamps.is_empty() || start > end {
1735 return vec![];
1736 }
1737
1738 (start..=end)
1739 .step_by(interval as usize)
1740 .map(|curr| {
1741 let mut offset = None;
1742 let mut length = 0;
1743 for (index, &ts) in timestamps.iter().enumerate() {
1744 if ts > curr - range && ts <= curr {
1745 offset.get_or_insert(index);
1746 length += 1;
1747 }
1748 }
1749 (offset.unwrap_or(0) as u32, length)
1750 })
1751 .collect()
1752 }
1753
1754 #[test]
1755 fn calculate_range_characterizes_returned_bounds() {
1756 let cases = [
1757 (
1758 "positive non-aligned last timestamp plus range",
1759 0,
1760 100,
1761 10,
1762 9,
1763 vec![13, 26],
1764 (20, 30),
1765 ),
1766 (
1767 "negative non-aligned last timestamp plus range",
1768 -50,
1769 50,
1770 10,
1771 9,
1772 vec![-37, -26],
1773 (-30, -20),
1774 ),
1775 (
1776 "query alignment not based on epoch",
1777 4,
1778 94,
1779 30,
1780 15,
1781 vec![20, 50, 80],
1782 (34, 94),
1783 ),
1784 ("leading data", 0, 100, 10, 10, vec![-10, 15], (0, 20)),
1785 (
1786 "trailing data past query end",
1787 0,
1788 100,
1789 10,
1790 10,
1791 vec![35, 45, 110],
1792 (40, 100),
1793 ),
1794 (
1795 "optimized start after query end",
1796 0,
1797 50,
1798 10,
1799 0,
1800 vec![100],
1801 (0, 50),
1802 ),
1803 ];
1804
1805 for (name, query_start, query_end, interval, range, timestamps, expected_bounds) in cases {
1806 let (_, bounds) =
1807 calculate_range_for_test(query_start, query_end, interval, range, ×tamps);
1808 assert_eq!(bounds, expected_bounds, "{name}");
1809 }
1810 }
1811
1812 #[test]
1813 fn calculate_range_keeps_extreme_range_tail() {
1814 let (ranges, bounds) =
1815 calculate_range_for_test(i64::MAX - 1, i64::MAX, 1, i64::MAX, &[i64::MAX]);
1816
1817 assert_eq!(bounds, (i64::MAX, i64::MAX));
1818 assert_eq!(ranges, vec![(0, 1)]);
1819 }
1820
1821 #[test]
1822 fn calculate_range_matches_bruteforce_oracle_for_deterministic_cases() {
1823 let cases = vec![
1824 (
1825 "duplicate lower and upper bounds",
1826 vec![0, 10, 10, 20, 20, 30],
1827 10,
1828 20,
1829 10,
1830 10,
1831 ),
1832 (
1833 "zero range excludes duplicates at current timestamp",
1834 vec![10, 10, 10],
1835 10,
1836 10,
1837 1,
1838 0,
1839 ),
1840 (
1841 "consecutive nonempty empty nonempty ranges",
1842 vec![10, 30],
1843 10,
1844 30,
1845 10,
1846 5,
1847 ),
1848 ("step smaller than range", vec![0, 4, 8, 12], 0, 12, 3, 5),
1849 ("step equal to range", vec![0, 5, 10, 15], 0, 15, 5, 5),
1850 ("step greater than range", vec![0, 7, 14, 21], 0, 21, 7, 3),
1851 (
1852 "negative sparse/tail timestamps",
1853 vec![-30, -20, -10, 0],
1854 -25,
1855 5,
1856 5,
1857 7,
1858 ),
1859 ("one sample", vec![42], 0, 100, 10, 15),
1860 ("empty input", vec![], -20, 20, 5, 10),
1861 ];
1862
1863 for (name, timestamps, query_start, query_end, interval, range) in cases {
1864 let (actual, (start, end)) =
1865 calculate_range_for_test(query_start, query_end, interval, range, ×tamps);
1866 let expected = calculate_range_oracle(×tamps, start, end, interval, range);
1867 assert_eq!(actual, expected, "{name}");
1868 }
1869 }
1870
1871 #[test]
1872 fn calculate_range_positive_time_translated_regression() {
1873 let timestamps = [0, 10, 20, 30];
1874 let expected = vec![(0, 1), (1, 1), (1, 1), (2, 1), (2, 1), (3, 1), (3, 1)];
1875 let (actual, bounds) = calculate_range_for_test(5, 35, 5, 7, ×tamps);
1876
1877 assert_eq!(bounds, (5, 35));
1878 assert_eq!(
1879 calculate_range_oracle(×tamps, bounds.0, bounds.1, 5, 7),
1880 expected
1881 );
1882 assert_eq!(actual, expected);
1883 }
1884
1885 #[test]
1886 fn calculate_range_matches_oracle_for_dense_positive_time_windows() {
1887 let timestamps = (0..=3_600).step_by(15).collect::<Vec<i64>>();
1888
1889 for (name, range) in [
1890 ("one minute", 60),
1891 ("five minutes", 300),
1892 ("one hour", 3_600),
1893 ] {
1894 let (actual, (start, end)) = calculate_range_for_test(0, 3_600, 15, range, ×tamps);
1895 assert_eq!((start, end), (0, 3_600), "{name} bounds");
1896 assert_eq!(
1897 actual,
1898 calculate_range_oracle(×tamps, start, end, 15, range),
1899 "{name} window"
1900 );
1901 }
1902 }
1903
1904 #[derive(Clone, Copy)]
1905 struct TinyPrng(u64);
1906
1907 impl TinyPrng {
1908 fn next_u64(&mut self) -> u64 {
1909 self.0 ^= self.0 << 13;
1910 self.0 ^= self.0 >> 7;
1911 self.0 ^= self.0 << 17;
1912 self.0
1913 }
1914
1915 fn next_i64(&mut self, min: i64, max: i64) -> i64 {
1916 min + (self.next_u64() % (max - min + 1) as u64) as i64
1917 }
1918 }
1919
1920 #[test]
1921 fn calculate_range_matches_bruteforce_oracle_for_seeded_matrix() {
1922 let mut prng = TinyPrng(0x5eed_cafe_f00d_baad);
1923
1924 for case in 0..512 {
1925 let interval = prng.next_i64(1, 11);
1926 let range = prng.next_i64(0, 25);
1927 let query_start = prng.next_i64(-200, 200);
1928 let query_end = query_start + interval * prng.next_i64(0, 20);
1929 let mut timestamps = Vec::new();
1930 let mut timestamp = prng.next_i64(-250, 250);
1931 for _ in 0..prng.next_i64(0, 20) {
1932 timestamp += prng.next_i64(0, 7);
1933 timestamps.push(timestamp);
1934 }
1935
1936 let (actual, (start, end)) =
1937 calculate_range_for_test(query_start, query_end, interval, range, ×tamps);
1938 let expected = calculate_range_oracle(×tamps, start, end, interval, range);
1939 let expected = if actual.is_empty() && !expected.is_empty() {
1940 assert!(
1941 expected.iter().all(|(_, len)| *len == 0),
1942 "case={case}, timestamps={timestamps:?}, query=({query_start}, {query_end}), \
1943 interval={interval}, range={range}, bounds=({start}, {end}): \
1944 no-intersection output must have no selected samples"
1945 );
1946 vec![]
1947 } else {
1948 expected
1949 };
1950 assert_eq!(
1951 actual, expected,
1952 "case={case}, timestamps={timestamps:?}, query=({query_start}, {query_end}), \
1953 interval={interval}, range={range}, bounds=({start}, {end})"
1954 );
1955 }
1956 }
1957}