1use std::cmp::Ordering;
16use std::collections::HashMap;
17use std::pin::Pin;
18use std::sync::Arc;
19use std::task::{Context, Poll};
20
21use datafusion::arrow::array::Array;
22use datafusion::common::tree_node::TreeNodeRecursion;
23use datafusion::common::{DFSchemaRef, Result as DataFusionResult};
24use datafusion::execution::context::TaskContext;
25use datafusion::logical_expr::{Expr, LogicalPlan, UserDefinedLogicalNodeCore};
26use datafusion::physical_expr::{
27 EquivalenceProperties, LexRequirement, OrderingRequirements, PhysicalSortRequirement,
28};
29use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
30use datafusion::physical_plan::expressions::Column as ColumnExpr;
31use datafusion::physical_plan::metrics::{BaselineMetrics, ExecutionPlanMetricsSet, MetricsSet};
32use datafusion::physical_plan::{
33 DisplayAs, DisplayFormatType, Distribution, ExecutionPlan, InputDistributionRequirements,
34 Partitioning, PhysicalExpr, PlanProperties, RecordBatchStream, SendableRecordBatchStream,
35};
36use datafusion_common::DFSchema;
37use datafusion_expr::{EmptyRelation, ident};
38use datatypes::arrow;
39use datatypes::arrow::array::{ArrayRef, Float64Array, TimestampMillisecondArray};
40use datatypes::arrow::datatypes::{DataType, Field, SchemaRef, TimeUnit};
41use datatypes::arrow::record_batch::RecordBatch;
42use datatypes::arrow_array::StringArray;
43use datatypes::compute::SortOptions;
44use futures::{Stream, StreamExt, ready};
45use greptime_proto::substrait_extension as pb;
46use prost::Message;
47use snafu::ResultExt;
48
49use crate::error::DeserializeSnafu;
50use crate::extension_plan::{Millisecond, resolve_column_name, serialize_column_index};
51
52#[derive(Debug, PartialEq, Eq, Hash)]
53pub struct Absent {
54 start: Millisecond,
55 end: Millisecond,
56 step: Millisecond,
57 time_index_column: String,
58 value_column: String,
59 fake_labels: Vec<(String, String)>,
60 input: LogicalPlan,
61 output_schema: DFSchemaRef,
62 unfix: Option<UnfixIndices>,
63}
64
65#[derive(Debug, PartialEq, Eq, Hash, PartialOrd)]
66struct UnfixIndices {
67 pub time_index_column_idx: u64,
68 pub value_column_idx: u64,
69}
70
71impl PartialOrd for Absent {
72 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
73 (
75 self.start,
76 self.end,
77 self.step,
78 &self.time_index_column,
79 &self.value_column,
80 &self.fake_labels,
81 )
82 .partial_cmp(&(
83 other.start,
84 other.end,
85 other.step,
86 &other.time_index_column,
87 &other.value_column,
88 &other.fake_labels,
89 ))
90 }
91}
92
93impl UserDefinedLogicalNodeCore for Absent {
94 fn name(&self) -> &str {
95 Self::name()
96 }
97
98 fn inputs(&self) -> Vec<&LogicalPlan> {
99 vec![&self.input]
100 }
101
102 fn schema(&self) -> &DFSchemaRef {
103 &self.output_schema
104 }
105
106 fn expressions(&self) -> Vec<Expr> {
107 if self.unfix.is_some() {
108 return vec![];
109 }
110
111 vec![ident(&self.time_index_column)]
112 }
113
114 fn necessary_children_exprs(&self, _output_columns: &[usize]) -> Option<Vec<Vec<usize>>> {
115 if self.unfix.is_some() {
116 return None;
117 }
118
119 let input_schema = self.input.schema();
120 let time_index_idx = input_schema.index_of_column_by_name(None, &self.time_index_column)?;
121 Some(vec![vec![time_index_idx]])
122 }
123
124 fn fmt_for_explain(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
125 write!(
126 f,
127 "PromAbsent: start={}, end={}, step={}",
128 self.start, self.end, self.step
129 )
130 }
131
132 fn with_exprs_and_inputs(
133 &self,
134 _exprs: Vec<Expr>,
135 inputs: Vec<LogicalPlan>,
136 ) -> DataFusionResult<Self> {
137 if inputs.is_empty() {
138 return Err(datafusion::error::DataFusionError::Internal(
139 "Absent must have at least one input".to_string(),
140 ));
141 }
142
143 let input: LogicalPlan = inputs[0].clone();
144 let input_schema = input.schema();
145
146 if let Some(unfix) = &self.unfix {
147 let time_index_column = resolve_column_name(
149 unfix.time_index_column_idx,
150 input_schema,
151 "Absent",
152 "time index",
153 )?;
154
155 let value_column =
156 resolve_column_name(unfix.value_column_idx, input_schema, "Absent", "value")?;
157
158 Self::try_new(
160 self.start,
161 self.end,
162 self.step,
163 time_index_column,
164 value_column,
165 self.fake_labels.clone(),
166 input,
167 )
168 } else {
169 Ok(Self {
170 start: self.start,
171 end: self.end,
172 step: self.step,
173 time_index_column: self.time_index_column.clone(),
174 value_column: self.value_column.clone(),
175 fake_labels: self.fake_labels.clone(),
176 input,
177 output_schema: self.output_schema.clone(),
178 unfix: None,
179 })
180 }
181 }
182}
183
184impl Absent {
185 pub fn try_new(
186 start: Millisecond,
187 end: Millisecond,
188 step: Millisecond,
189 time_index_column: String,
190 value_column: String,
191 fake_labels: Vec<(String, String)>,
192 input: LogicalPlan,
193 ) -> DataFusionResult<Self> {
194 let mut fields = vec![
195 Field::new(
196 &time_index_column,
197 DataType::Timestamp(TimeUnit::Millisecond, None),
198 true,
199 ),
200 Field::new(&value_column, DataType::Float64, true),
201 ];
202
203 let mut fake_labels = fake_labels
205 .into_iter()
206 .collect::<HashMap<String, String>>()
207 .into_iter()
208 .collect::<Vec<_>>();
209 fake_labels.sort_unstable_by(|a, b| a.0.cmp(&b.0));
210 for (name, _) in fake_labels.iter() {
211 fields.push(Field::new(name, DataType::Utf8, true));
212 }
213
214 let output_schema = Arc::new(DFSchema::from_unqualified_fields(
215 fields.into(),
216 HashMap::new(),
217 )?);
218
219 Ok(Self {
220 start,
221 end,
222 step,
223 time_index_column,
224 value_column,
225 fake_labels,
226 input,
227 output_schema,
228 unfix: None,
229 })
230 }
231
232 pub const fn name() -> &'static str {
233 "prom_absent"
234 }
235
236 pub fn to_execution_plan(&self, exec_input: Arc<dyn ExecutionPlan>) -> Arc<dyn ExecutionPlan> {
237 let output_schema = Arc::new(self.output_schema.as_arrow().clone());
238 let properties = Arc::new(PlanProperties::new(
239 EquivalenceProperties::new(output_schema.clone()),
240 Partitioning::UnknownPartitioning(1),
241 EmissionType::Incremental,
242 Boundedness::Bounded,
243 ));
244 Arc::new(AbsentExec {
245 start: self.start,
246 end: self.end,
247 step: self.step,
248 time_index_column: self.time_index_column.clone(),
249 value_column: self.value_column.clone(),
250 fake_labels: self.fake_labels.clone(),
251 output_schema: output_schema.clone(),
252 input: exec_input,
253 properties,
254 metric: ExecutionPlanMetricsSet::new(),
255 })
256 }
257
258 pub fn serialize(&self) -> Vec<u8> {
259 let time_index_column_idx =
260 serialize_column_index(self.input.schema(), &self.time_index_column);
261
262 let value_column_idx = serialize_column_index(self.input.schema(), &self.value_column);
263
264 pb::Absent {
265 start: self.start,
266 end: self.end,
267 step: self.step,
268 time_index_column_idx,
269 value_column_idx,
270 fake_labels: self
271 .fake_labels
272 .iter()
273 .map(|(name, value)| pb::LabelPair {
274 key: name.clone(),
275 value: value.clone(),
276 })
277 .collect(),
278 ..Default::default()
279 }
280 .encode_to_vec()
281 }
282
283 pub fn deserialize(bytes: &[u8]) -> DataFusionResult<Self> {
284 let pb_absent = pb::Absent::decode(bytes).context(DeserializeSnafu)?;
285 let placeholder_plan = LogicalPlan::EmptyRelation(EmptyRelation {
286 produce_one_row: false,
287 schema: Arc::new(DFSchema::empty()),
288 });
289
290 let unfix = UnfixIndices {
291 time_index_column_idx: pb_absent.time_index_column_idx,
292 value_column_idx: pb_absent.value_column_idx,
293 };
294
295 Ok(Self {
296 start: pb_absent.start,
297 end: pb_absent.end,
298 step: pb_absent.step,
299 time_index_column: String::new(),
300 value_column: String::new(),
301 fake_labels: pb_absent
302 .fake_labels
303 .iter()
304 .map(|label| (label.key.clone(), label.value.clone()))
305 .collect(),
306 input: placeholder_plan,
307 output_schema: Arc::new(DFSchema::empty()),
308 unfix: Some(unfix),
309 })
310 }
311}
312
313#[derive(Debug)]
314pub struct AbsentExec {
315 start: Millisecond,
316 end: Millisecond,
317 step: Millisecond,
318 time_index_column: String,
319 value_column: String,
320 fake_labels: Vec<(String, String)>,
321 output_schema: SchemaRef,
322 input: Arc<dyn ExecutionPlan>,
323 properties: Arc<PlanProperties>,
324 metric: ExecutionPlanMetricsSet,
325}
326
327impl ExecutionPlan for AbsentExec {
328 fn apply_expressions(
329 &self,
330 _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> datafusion_common::Result<TreeNodeRecursion>,
331 ) -> DataFusionResult<TreeNodeRecursion> {
332 Ok(TreeNodeRecursion::Continue)
333 }
334
335 fn schema(&self) -> SchemaRef {
336 self.output_schema.clone()
337 }
338
339 fn properties(&self) -> &Arc<PlanProperties> {
340 &self.properties
341 }
342
343 fn input_distribution_requirements(&self) -> InputDistributionRequirements {
344 InputDistributionRequirements::new(vec![Distribution::SinglePartition])
345 }
346
347 fn required_input_ordering(&self) -> Vec<Option<OrderingRequirements>> {
348 let requirement = LexRequirement::from([PhysicalSortRequirement {
349 expr: Arc::new(
350 ColumnExpr::new_with_schema(&self.time_index_column, &self.input.schema()).unwrap(),
351 ),
352 options: Some(SortOptions {
353 descending: false,
354 nulls_first: false,
355 }),
356 }]);
357 vec![Some(OrderingRequirements::new(requirement))]
358 }
359
360 fn maintains_input_order(&self) -> Vec<bool> {
361 vec![false]
362 }
363
364 fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
365 vec![&self.input]
366 }
367
368 fn with_new_children(
369 self: Arc<Self>,
370 children: Vec<Arc<dyn ExecutionPlan>>,
371 ) -> DataFusionResult<Arc<dyn ExecutionPlan>> {
372 assert!(!children.is_empty());
373 Ok(Arc::new(Self {
374 start: self.start,
375 end: self.end,
376 step: self.step,
377 time_index_column: self.time_index_column.clone(),
378 value_column: self.value_column.clone(),
379 fake_labels: self.fake_labels.clone(),
380 output_schema: self.output_schema.clone(),
381 input: children[0].clone(),
382 properties: self.properties.clone(),
383 metric: self.metric.clone(),
384 }))
385 }
386
387 fn execute(
388 &self,
389 partition: usize,
390 context: Arc<TaskContext>,
391 ) -> DataFusionResult<SendableRecordBatchStream> {
392 let baseline_metric = BaselineMetrics::new(&self.metric, partition);
393 let batch_size = context.session_config().batch_size();
394 let input = self.input.execute(partition, context)?;
395
396 Ok(Box::pin(AbsentStream {
397 end: self.end,
398 step: self.step,
399 batch_size,
400 time_index_column_index: self
401 .input
402 .schema()
403 .column_with_name(&self.time_index_column)
404 .unwrap() .0,
406 output_schema: self.output_schema.clone(),
407 fake_labels: self.fake_labels.clone(),
408 input,
409 metric: baseline_metric,
410 output_timestamps: Vec::new(),
412 input_timestamps: Vec::new(),
413 input_timestamp_offset: 0,
414 output_ts_cursor: self.start,
416 input_finished: false,
417 }))
418 }
419
420 fn metrics(&self) -> Option<MetricsSet> {
421 Some(self.metric.clone_inner())
422 }
423
424 fn name(&self) -> &str {
425 "AbsentExec"
426 }
427}
428
429impl DisplayAs for AbsentExec {
430 fn fmt_as(&self, t: DisplayFormatType, f: &mut std::fmt::Formatter) -> std::fmt::Result {
431 match t {
432 DisplayFormatType::Default
433 | DisplayFormatType::Verbose
434 | DisplayFormatType::TreeRender => {
435 write!(
436 f,
437 "PromAbsentExec: start={}, end={}, step={}",
438 self.start, self.end, self.step
439 )
440 }
441 }
442 }
443}
444
445pub struct AbsentStream {
446 end: Millisecond,
447 step: Millisecond,
448 batch_size: usize,
449 time_index_column_index: usize,
450 output_schema: SchemaRef,
451 fake_labels: Vec<(String, String)>,
452 input: SendableRecordBatchStream,
453 metric: BaselineMetrics,
454 output_timestamps: Vec<Millisecond>,
456 input_timestamps: Vec<Millisecond>,
458 input_timestamp_offset: usize,
459 output_ts_cursor: Millisecond,
461 input_finished: bool,
462}
463
464impl RecordBatchStream for AbsentStream {
465 fn schema(&self) -> SchemaRef {
466 self.output_schema.clone()
467 }
468}
469
470impl Stream for AbsentStream {
471 type Item = DataFusionResult<RecordBatch>;
472
473 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
474 loop {
475 if self.has_pending_input_timestamps() {
476 let timer = std::time::Instant::now();
477 if let Err(e) = self.process_input_batch() {
478 return Poll::Ready(Some(Err(e)));
479 }
480 self.metric.elapsed_compute().add_elapsed(timer);
481
482 match self.flush_output_batch() {
483 Ok(Some(batch)) => return Poll::Ready(Some(Ok(batch))),
484 Ok(None) => continue,
485 Err(e) => return Poll::Ready(Some(Err(e))),
486 }
487 }
488
489 if self.input_finished {
490 let timer = std::time::Instant::now();
491 if let Err(e) = self.process_remaining_absent_timestamps() {
492 return Poll::Ready(Some(Err(e)));
493 }
494 self.metric.elapsed_compute().add_elapsed(timer);
495
496 match self.flush_output_batch() {
497 Ok(Some(batch)) => return Poll::Ready(Some(Ok(batch))),
498 Ok(None) => return Poll::Ready(None),
499 Err(e) => return Poll::Ready(Some(Err(e))),
500 }
501 }
502
503 match ready!(self.input.poll_next_unpin(cx)) {
504 Some(Ok(batch)) => {
505 let timer = std::time::Instant::now();
506 if let Err(e) = self.buffer_input_timestamps(&batch) {
507 return Poll::Ready(Some(Err(e)));
508 }
509 self.metric.elapsed_compute().add_elapsed(timer);
510 }
511 Some(Err(e)) => return Poll::Ready(Some(Err(e))),
512 None => {
513 self.input_finished = true;
514 }
515 }
516 }
517 }
518}
519
520impl AbsentStream {
521 fn buffer_input_timestamps(&mut self, batch: &RecordBatch) -> DataFusionResult<()> {
522 let timestamp_array = batch.column(self.time_index_column_index);
523 let milli_ts_array = arrow::compute::cast(
524 timestamp_array,
525 &DataType::Timestamp(TimeUnit::Millisecond, None),
526 )?;
527 let timestamp_array = milli_ts_array
528 .as_any()
529 .downcast_ref::<TimestampMillisecondArray>()
530 .unwrap();
531 self.input_timestamps.clear();
532 self.input_timestamps
533 .extend_from_slice(timestamp_array.values());
534 self.input_timestamp_offset = 0;
535 Ok(())
536 }
537
538 fn has_pending_input_timestamps(&self) -> bool {
539 self.input_timestamp_offset < self.input_timestamps.len()
540 }
541
542 fn process_input_batch(&mut self) -> DataFusionResult<()> {
543 while self.input_timestamp_offset < self.input_timestamps.len() {
544 let input_ts = self.input_timestamps[self.input_timestamp_offset];
545
546 while self.output_ts_cursor < input_ts && self.output_ts_cursor <= self.end {
548 self.output_timestamps.push(self.output_ts_cursor);
549 self.output_ts_cursor += self.step;
550
551 if self.output_timestamps.len() >= self.batch_size {
552 return Ok(());
553 }
554 }
555
556 if self.output_ts_cursor == input_ts {
558 self.output_ts_cursor += self.step;
559 }
560
561 self.input_timestamp_offset += 1;
562 }
563
564 self.input_timestamps.clear();
565 self.input_timestamp_offset = 0;
566 Ok(())
567 }
568
569 fn process_remaining_absent_timestamps(&mut self) -> DataFusionResult<()> {
570 while self.output_ts_cursor <= self.end {
571 self.output_timestamps.push(self.output_ts_cursor);
572 self.output_ts_cursor += self.step;
573
574 if self.output_timestamps.len() >= self.batch_size {
575 return Ok(());
576 }
577 }
578 Ok(())
579 }
580
581 fn flush_output_batch(&mut self) -> DataFusionResult<Option<RecordBatch>> {
582 if self.output_timestamps.is_empty() {
583 return Ok(None);
584 }
585
586 let timestamps = if self.output_timestamps.len() <= self.batch_size {
587 std::mem::take(&mut self.output_timestamps)
588 } else {
589 let remaining = self.output_timestamps.split_off(self.batch_size);
590 std::mem::replace(&mut self.output_timestamps, remaining)
591 };
592
593 let mut columns: Vec<ArrayRef> = Vec::with_capacity(self.output_schema.fields().len());
594 let num_rows = timestamps.len();
595 columns.push(Arc::new(TimestampMillisecondArray::from(timestamps)) as _);
596 columns.push(Arc::new(Float64Array::from(vec![1.0; num_rows])) as _);
597
598 for (_, value) in self.fake_labels.iter() {
599 columns.push(Arc::new(StringArray::from_iter(std::iter::repeat_n(
600 Some(value.clone()),
601 num_rows,
602 ))) as _);
603 }
604
605 let batch = RecordBatch::try_new(self.output_schema.clone(), columns)?;
606
607 Ok(Some(batch))
608 }
609}
610
611#[cfg(test)]
612mod tests {
613 use std::sync::Arc;
614
615 use datafusion::arrow::datatypes::{DataType, Field, Schema, TimeUnit};
616 use datafusion::arrow::record_batch::RecordBatch;
617 use datafusion::catalog::memory::DataSourceExec;
618 use datafusion::datasource::memory::MemorySourceConfig;
619 use datafusion::prelude::{SessionConfig, SessionContext};
620 use datatypes::arrow::array::{Float64Array, TimestampMillisecondArray};
621
622 use super::*;
623
624 #[tokio::test]
625 async fn test_absent_basic() {
626 let schema = Arc::new(Schema::new(vec![
627 Field::new(
628 "timestamp",
629 DataType::Timestamp(TimeUnit::Millisecond, None),
630 true,
631 ),
632 Field::new("value", DataType::Float64, true),
633 ]));
634
635 let timestamp_array = Arc::new(TimestampMillisecondArray::from(vec![0, 2000, 4000]));
637 let value_array = Arc::new(Float64Array::from(vec![1.0, 2.0, 3.0]));
638 let batch =
639 RecordBatch::try_new(schema.clone(), vec![timestamp_array, value_array]).unwrap();
640
641 let memory_exec = DataSourceExec::new(Arc::new(
642 MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(),
643 ));
644
645 let output_schema = Arc::new(Schema::new(vec![
646 Field::new(
647 "timestamp",
648 DataType::Timestamp(TimeUnit::Millisecond, None),
649 true,
650 ),
651 Field::new("value", DataType::Float64, true),
652 ]));
653
654 let absent_exec = AbsentExec {
655 start: 0,
656 end: 5000,
657 step: 1000,
658 time_index_column: "timestamp".to_string(),
659 value_column: "value".to_string(),
660 fake_labels: vec![],
661 output_schema: output_schema.clone(),
662 input: Arc::new(memory_exec),
663 properties: Arc::new(PlanProperties::new(
664 EquivalenceProperties::new(output_schema.clone()),
665 Partitioning::UnknownPartitioning(1),
666 EmissionType::Incremental,
667 Boundedness::Bounded,
668 )),
669 metric: ExecutionPlanMetricsSet::new(),
670 };
671
672 let session_ctx = SessionContext::new();
673 let task_ctx = session_ctx.task_ctx();
674 let mut stream = absent_exec.execute(0, task_ctx).unwrap();
675
676 let mut output_timestamps = Vec::new();
678 while let Some(batch_result) = stream.next().await {
679 let batch = batch_result.unwrap();
680 let ts_array = batch
681 .column(0)
682 .as_any()
683 .downcast_ref::<TimestampMillisecondArray>()
684 .unwrap();
685 for i in 0..ts_array.len() {
686 if !ts_array.is_null(i) {
687 let ts = ts_array.value(i);
688 output_timestamps.push(ts);
689 }
690 }
691 }
692
693 assert_eq!(output_timestamps, vec![1000, 3000, 5000]);
696 }
697
698 #[tokio::test]
699 async fn test_absent_empty_input() {
700 let schema = Arc::new(Schema::new(vec![
701 Field::new(
702 "timestamp",
703 DataType::Timestamp(TimeUnit::Millisecond, None),
704 true,
705 ),
706 Field::new("value", DataType::Float64, true),
707 ]));
708
709 let memory_exec = DataSourceExec::new(Arc::new(
711 MemorySourceConfig::try_new(&[vec![]], schema, None).unwrap(),
712 ));
713
714 let output_schema = Arc::new(Schema::new(vec![
715 Field::new(
716 "timestamp",
717 DataType::Timestamp(TimeUnit::Millisecond, None),
718 true,
719 ),
720 Field::new("value", DataType::Float64, true),
721 ]));
722 let absent_exec = AbsentExec {
723 start: 0,
724 end: 2000,
725 step: 1000,
726 time_index_column: "timestamp".to_string(),
727 value_column: "value".to_string(),
728 fake_labels: vec![],
729 output_schema: output_schema.clone(),
730 input: Arc::new(memory_exec),
731 properties: Arc::new(PlanProperties::new(
732 EquivalenceProperties::new(output_schema.clone()),
733 Partitioning::UnknownPartitioning(1),
734 EmissionType::Incremental,
735 Boundedness::Bounded,
736 )),
737 metric: ExecutionPlanMetricsSet::new(),
738 };
739
740 let session_ctx = SessionContext::new();
741 let task_ctx = session_ctx.task_ctx();
742 let mut stream = absent_exec.execute(0, task_ctx).unwrap();
743
744 let mut output_timestamps = Vec::new();
746 while let Some(batch_result) = stream.next().await {
747 let batch = batch_result.unwrap();
748 let ts_array = batch
749 .column(0)
750 .as_any()
751 .downcast_ref::<TimestampMillisecondArray>()
752 .unwrap();
753 for i in 0..ts_array.len() {
754 if !ts_array.is_null(i) {
755 let ts = ts_array.value(i);
756 output_timestamps.push(ts);
757 }
758 }
759 }
760
761 assert_eq!(output_timestamps, vec![0, 1000, 2000]);
763 }
764
765 #[tokio::test]
766 async fn test_absent_respects_session_batch_size_for_large_gap() {
767 let schema = Arc::new(Schema::new(vec![
768 Field::new(
769 "timestamp",
770 DataType::Timestamp(TimeUnit::Millisecond, None),
771 true,
772 ),
773 Field::new("value", DataType::Float64, true),
774 ]));
775
776 let timestamp_array = Arc::new(TimestampMillisecondArray::from(vec![9]));
777 let value_array = Arc::new(Float64Array::from(vec![1.0]));
778 let batch =
779 RecordBatch::try_new(schema.clone(), vec![timestamp_array, value_array]).unwrap();
780
781 let memory_exec = DataSourceExec::new(Arc::new(
782 MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(),
783 ));
784
785 let output_schema = Arc::new(Schema::new(vec![
786 Field::new(
787 "timestamp",
788 DataType::Timestamp(TimeUnit::Millisecond, None),
789 true,
790 ),
791 Field::new("value", DataType::Float64, true),
792 ]));
793
794 let absent_exec = AbsentExec {
795 start: 0,
796 end: 10,
797 step: 1,
798 time_index_column: "timestamp".to_string(),
799 value_column: "value".to_string(),
800 fake_labels: vec![],
801 output_schema: output_schema.clone(),
802 input: Arc::new(memory_exec),
803 properties: Arc::new(PlanProperties::new(
804 EquivalenceProperties::new(output_schema.clone()),
805 Partitioning::UnknownPartitioning(1),
806 EmissionType::Incremental,
807 Boundedness::Bounded,
808 )),
809 metric: ExecutionPlanMetricsSet::new(),
810 };
811
812 let session_ctx = SessionContext::new_with_config(SessionConfig::new().with_batch_size(3));
813 let task_ctx = session_ctx.task_ctx();
814 let mut stream = absent_exec.execute(0, task_ctx).unwrap();
815
816 let mut batch_sizes = Vec::new();
817 let mut output_timestamps = Vec::new();
818 while let Some(batch_result) = stream.next().await {
819 let batch = batch_result.unwrap();
820 batch_sizes.push(batch.num_rows());
821
822 let ts_array = batch
823 .column(0)
824 .as_any()
825 .downcast_ref::<TimestampMillisecondArray>()
826 .unwrap();
827 for i in 0..ts_array.len() {
828 if !ts_array.is_null(i) {
829 output_timestamps.push(ts_array.value(i));
830 }
831 }
832 }
833
834 assert_eq!(batch_sizes, vec![3, 3, 3, 1]);
835 assert_eq!(output_timestamps, vec![0, 1, 2, 3, 4, 5, 6, 7, 8, 10]);
836 }
837
838 #[tokio::test]
839 async fn test_absent_resumes_same_input_timestamp_after_batch_flush() {
840 let schema = Arc::new(Schema::new(vec![
841 Field::new(
842 "timestamp",
843 DataType::Timestamp(TimeUnit::Millisecond, None),
844 true,
845 ),
846 Field::new("value", DataType::Float64, true),
847 ]));
848
849 let timestamp_array = Arc::new(TimestampMillisecondArray::from(vec![9]));
850 let value_array = Arc::new(Float64Array::from(vec![1.0]));
851 let batch =
852 RecordBatch::try_new(schema.clone(), vec![timestamp_array, value_array]).unwrap();
853
854 let memory_exec = DataSourceExec::new(Arc::new(
855 MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(),
856 ));
857
858 let output_schema = Arc::new(Schema::new(vec![
859 Field::new(
860 "timestamp",
861 DataType::Timestamp(TimeUnit::Millisecond, None),
862 true,
863 ),
864 Field::new("value", DataType::Float64, true),
865 ]));
866
867 let absent_exec = AbsentExec {
868 start: 0,
869 end: 9,
870 step: 1,
871 time_index_column: "timestamp".to_string(),
872 value_column: "value".to_string(),
873 fake_labels: vec![],
874 output_schema: output_schema.clone(),
875 input: Arc::new(memory_exec),
876 properties: Arc::new(PlanProperties::new(
877 EquivalenceProperties::new(output_schema.clone()),
878 Partitioning::UnknownPartitioning(1),
879 EmissionType::Incremental,
880 Boundedness::Bounded,
881 )),
882 metric: ExecutionPlanMetricsSet::new(),
883 };
884
885 let session_ctx = SessionContext::new_with_config(SessionConfig::new().with_batch_size(3));
886 let task_ctx = session_ctx.task_ctx();
887 let mut stream = absent_exec.execute(0, task_ctx).unwrap();
888
889 let mut output_timestamps = Vec::new();
890 while let Some(batch_result) = stream.next().await {
891 let batch = batch_result.unwrap();
892 let ts_array = batch
893 .column(0)
894 .as_any()
895 .downcast_ref::<TimestampMillisecondArray>()
896 .unwrap();
897 for i in 0..ts_array.len() {
898 if !ts_array.is_null(i) {
899 output_timestamps.push(ts_array.value(i));
900 }
901 }
902 }
903
904 assert_eq!(output_timestamps, vec![0, 1, 2, 3, 4, 5, 6, 7, 8]);
905 }
906}