Skip to main content

promql/extension_plan/
absent.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::cmp::Ordering;
16use std::collections::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        // compare on fields except schema and input
74        (
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            // transform indices to names
148            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            // Recreate output schema with actual field names
159            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        // remove duplicate fake labels
204        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() // Safety: we have checked the column name in `try_new`
405                .0,
406            output_schema: self.output_schema.clone(),
407            fake_labels: self.fake_labels.clone(),
408            input,
409            metric: baseline_metric,
410            // Buffer for streaming output timestamps
411            output_timestamps: Vec::new(),
412            input_timestamps: Vec::new(),
413            input_timestamp_offset: 0,
414            // Current timestamp in the output range
415            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    // Buffer for streaming output timestamps
455    output_timestamps: Vec<Millisecond>,
456    // Current input timestamps being processed incrementally.
457    input_timestamps: Vec<Millisecond>,
458    input_timestamp_offset: usize,
459    // Current timestamp in the output range
460    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            // Generate absent timestamps up to this input timestamp
547            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            // Skip the input timestamp if it matches our cursor
557            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        // Input has timestamps: 0, 2000, 4000
636        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        // Collect all output batches
677        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        // Should output absent timestamps: 1000, 3000, 5000
694        // (0, 2000, 4000 exist in input, so 1000, 3000, 5000 are absent)
695        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        // Empty input
710        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        // Collect all output timestamps
745        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        // Should output all timestamps in range: 0, 1000, 2000
762        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}