Skip to main content

promql/extension_plan/
normalize.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::any::Any;
16use std::pin::Pin;
17use std::sync::Arc;
18use std::task::{Context, Poll};
19
20use common_query::native_histogram::{START_TIMESTAMP_FIELD, native_histogram_arrow_type};
21use datafusion::arrow::array::{Array, BooleanArray, StructArray};
22use datafusion::arrow::compute;
23use datafusion::common::{DFSchema, DFSchemaRef, Result as DataFusionResult, Statistics};
24use datafusion::error::DataFusionError;
25use datafusion::execution::context::TaskContext;
26use datafusion::logical_expr::{EmptyRelation, Expr, LogicalPlan, UserDefinedLogicalNodeCore};
27use datafusion::physical_plan::expressions::Column as ColumnExpr;
28use datafusion::physical_plan::metrics::{
29    BaselineMetrics, Count, ExecutionPlanMetricsSet, MetricBuilder, MetricValue, MetricsSet,
30};
31use datafusion::physical_plan::{
32    DisplayAs, DisplayFormatType, Distribution, ExecutionPlan, PlanProperties, RecordBatchStream,
33    SendableRecordBatchStream,
34};
35use datafusion_expr::col;
36use datatypes::arrow::array::TimestampMillisecondArray;
37use datatypes::arrow::datatypes::{SchemaRef, TimestampMillisecondType};
38use datatypes::arrow::record_batch::RecordBatch;
39use futures::{Stream, StreamExt, ready};
40use greptime_proto::substrait_extension as pb;
41use prost::Message;
42use snafu::ResultExt;
43
44use crate::error::{DeserializeSnafu, Result};
45use crate::extension_plan::{
46    METRIC_NUM_SERIES, Millisecond, is_prometheus_stale_sample, prometheus_stale_sample_column,
47    resolve_column_name, serialize_column_index,
48};
49use crate::metrics::PROMQL_SERIES_COUNT;
50
51/// Normalize the input record batch. Notice that for simplicity, this method assumes
52/// the input batch only contains sample points from one time series.
53///
54/// Roughly speaking, this method does these things:
55/// - bias sample and native histogram start timestamps by offset
56/// - sort the record batch based on timestamp column
57/// - remove Prometheus stale markers (optional)
58#[derive(Debug, PartialEq, Eq, Hash, PartialOrd)]
59pub struct SeriesNormalize {
60    offset: Millisecond,
61    time_index_column_name: String,
62    filter_stale_markers: bool,
63    tag_columns: Vec<String>,
64
65    input: LogicalPlan,
66    unfix: Option<UnfixIndices>,
67}
68
69#[derive(Debug, PartialEq, Eq, Hash, PartialOrd)]
70struct UnfixIndices {
71    pub time_index_idx: u64,
72    pub tag_column_indices: Vec<u64>,
73}
74
75impl UserDefinedLogicalNodeCore for SeriesNormalize {
76    fn name(&self) -> &str {
77        Self::name()
78    }
79
80    fn inputs(&self) -> Vec<&LogicalPlan> {
81        vec![&self.input]
82    }
83
84    fn schema(&self) -> &DFSchemaRef {
85        self.input.schema()
86    }
87
88    fn expressions(&self) -> Vec<datafusion::logical_expr::Expr> {
89        if self.unfix.is_some() {
90            return vec![];
91        }
92
93        self.tag_columns
94            .iter()
95            .map(col)
96            .chain(std::iter::once(col(&self.time_index_column_name)))
97            .collect()
98    }
99
100    fn necessary_children_exprs(&self, output_columns: &[usize]) -> Option<Vec<Vec<usize>>> {
101        if self.unfix.is_some() {
102            return None;
103        }
104
105        let input_schema = self.input.schema();
106        if output_columns.is_empty() {
107            let indices = (0..input_schema.fields().len()).collect::<Vec<_>>();
108            return Some(vec![indices]);
109        }
110
111        let mut required = Vec::with_capacity(output_columns.len() + 1 + self.tag_columns.len());
112        required.extend_from_slice(output_columns);
113        required.push(input_schema.index_of_column_by_name(None, &self.time_index_column_name)?);
114        for tag in &self.tag_columns {
115            required.push(input_schema.index_of_column_by_name(None, tag)?);
116        }
117
118        required.sort_unstable();
119        required.dedup();
120        Some(vec![required])
121    }
122
123    fn fmt_for_explain(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
124        write!(
125            f,
126            "PromSeriesNormalize: offset=[{}], time index=[{}], filter NaN: [{}]",
127            self.offset, self.time_index_column_name, self.filter_stale_markers
128        )
129    }
130
131    fn with_exprs_and_inputs(
132        &self,
133        _exprs: Vec<Expr>,
134        inputs: Vec<LogicalPlan>,
135    ) -> DataFusionResult<Self> {
136        if inputs.is_empty() {
137            return Err(DataFusionError::Internal(
138                "SeriesNormalize should have at least one input".to_string(),
139            ));
140        }
141
142        let input: LogicalPlan = inputs.into_iter().next().unwrap();
143        let input_schema = input.schema();
144
145        if let Some(unfix) = &self.unfix {
146            // transform indices to names
147            let time_index_column_name = resolve_column_name(
148                unfix.time_index_idx,
149                input_schema,
150                "SeriesNormalize",
151                "time index",
152            )?;
153
154            let tag_columns = unfix
155                .tag_column_indices
156                .iter()
157                .map(|idx| resolve_column_name(*idx, input_schema, "SeriesNormalize", "tag"))
158                .collect::<DataFusionResult<Vec<String>>>()?;
159
160            Ok(Self {
161                offset: self.offset,
162                time_index_column_name,
163                filter_stale_markers: self.filter_stale_markers,
164                tag_columns,
165                input,
166                unfix: None,
167            })
168        } else {
169            Ok(Self {
170                offset: self.offset,
171                time_index_column_name: self.time_index_column_name.clone(),
172                filter_stale_markers: self.filter_stale_markers,
173                tag_columns: self.tag_columns.clone(),
174                input,
175                unfix: None,
176            })
177        }
178    }
179}
180
181impl SeriesNormalize {
182    pub fn new<N: AsRef<str>>(
183        offset: Millisecond,
184        time_index_column_name: N,
185        filter_stale_markers: bool,
186        tag_columns: Vec<String>,
187        input: LogicalPlan,
188    ) -> Self {
189        Self {
190            offset,
191            time_index_column_name: time_index_column_name.as_ref().to_string(),
192            filter_stale_markers,
193            tag_columns,
194            input,
195            unfix: None,
196        }
197    }
198
199    pub const fn name() -> &'static str {
200        "SeriesNormalize"
201    }
202
203    pub fn to_execution_plan(&self, exec_input: Arc<dyn ExecutionPlan>) -> Arc<dyn ExecutionPlan> {
204        Arc::new(SeriesNormalizeExec {
205            offset: self.offset,
206            time_index_column_name: self.time_index_column_name.clone(),
207            filter_stale_markers: self.filter_stale_markers,
208            input: exec_input,
209            tag_columns: self.tag_columns.clone(),
210            metric: ExecutionPlanMetricsSet::new(),
211        })
212    }
213
214    pub fn serialize(&self) -> Vec<u8> {
215        let time_index_idx =
216            serialize_column_index(self.input.schema(), &self.time_index_column_name);
217
218        let tag_column_indices = self
219            .tag_columns
220            .iter()
221            .map(|name| serialize_column_index(self.input.schema(), name))
222            .collect::<Vec<u64>>();
223
224        pb::SeriesNormalize {
225            offset: self.offset,
226            time_index_idx,
227            filter_nan: self.filter_stale_markers,
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_normalize = pb::SeriesNormalize::decode(bytes).context(DeserializeSnafu)?;
236        let placeholder_plan = LogicalPlan::EmptyRelation(EmptyRelation {
237            produce_one_row: false,
238            schema: Arc::new(DFSchema::empty()),
239        });
240
241        let unfix = UnfixIndices {
242            time_index_idx: pb_normalize.time_index_idx,
243            tag_column_indices: pb_normalize.tag_column_indices.clone(),
244        };
245
246        Ok(Self {
247            offset: pb_normalize.offset,
248            time_index_column_name: String::new(),
249            filter_stale_markers: pb_normalize.filter_nan,
250            tag_columns: Vec::new(),
251            input: placeholder_plan,
252            unfix: Some(unfix),
253        })
254    }
255}
256
257#[derive(Debug)]
258pub struct SeriesNormalizeExec {
259    offset: Millisecond,
260    time_index_column_name: String,
261    filter_stale_markers: bool,
262    tag_columns: Vec<String>,
263
264    input: Arc<dyn ExecutionPlan>,
265    metric: ExecutionPlanMetricsSet,
266}
267
268impl ExecutionPlan for SeriesNormalizeExec {
269    fn as_any(&self) -> &dyn Any {
270        self
271    }
272
273    fn schema(&self) -> SchemaRef {
274        self.input.schema()
275    }
276
277    fn required_input_distribution(&self) -> Vec<Distribution> {
278        if self.tag_columns.is_empty() {
279            return vec![Distribution::SinglePartition];
280        }
281
282        let schema = self.input.schema();
283        vec![Distribution::HashPartitioned(
284            self.tag_columns
285                .iter()
286                // Safety: the tag column names is verified in the planning phase
287                .map(|tag| Arc::new(ColumnExpr::new_with_schema(tag, &schema).unwrap()) as _)
288                .collect(),
289        )]
290    }
291
292    fn properties(&self) -> &Arc<PlanProperties> {
293        self.input.properties()
294    }
295
296    fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
297        vec![&self.input]
298    }
299
300    fn with_new_children(
301        self: Arc<Self>,
302        children: Vec<Arc<dyn ExecutionPlan>>,
303    ) -> DataFusionResult<Arc<dyn ExecutionPlan>> {
304        assert!(!children.is_empty());
305        Ok(Arc::new(Self {
306            offset: self.offset,
307            time_index_column_name: self.time_index_column_name.clone(),
308            filter_stale_markers: self.filter_stale_markers,
309            input: children[0].clone(),
310            tag_columns: self.tag_columns.clone(),
311            metric: self.metric.clone(),
312        }))
313    }
314
315    fn execute(
316        &self,
317        partition: usize,
318        context: Arc<TaskContext>,
319    ) -> DataFusionResult<SendableRecordBatchStream> {
320        let baseline_metric = BaselineMetrics::new(&self.metric, partition);
321        let metrics_builder = MetricBuilder::new(&self.metric);
322        let num_series = Count::new();
323        metrics_builder
324            .with_partition(partition)
325            .build(MetricValue::Count {
326                name: METRIC_NUM_SERIES.into(),
327                count: num_series.clone(),
328            });
329
330        let input = self.input.execute(partition, context)?;
331        let schema = input.schema();
332        let time_index = schema
333            .column_with_name(&self.time_index_column_name)
334            .expect("time index column not found")
335            .0;
336        Ok(Box::pin(SeriesNormalizeStream {
337            offset: self.offset,
338            time_index,
339            filter_stale_markers: self.filter_stale_markers,
340            schema,
341            input,
342            metric: baseline_metric,
343            num_series,
344        }))
345    }
346
347    fn metrics(&self) -> Option<MetricsSet> {
348        Some(self.metric.clone_inner())
349    }
350
351    fn partition_statistics(&self, partition: Option<usize>) -> DataFusionResult<Statistics> {
352        self.input.partition_statistics(partition)
353    }
354
355    fn name(&self) -> &str {
356        "SeriesNormalizeExec"
357    }
358}
359
360impl DisplayAs for SeriesNormalizeExec {
361    fn fmt_as(&self, t: DisplayFormatType, f: &mut std::fmt::Formatter) -> std::fmt::Result {
362        match t {
363            DisplayFormatType::Default
364            | DisplayFormatType::Verbose
365            | DisplayFormatType::TreeRender => {
366                write!(
367                    f,
368                    "PromSeriesNormalizeExec: offset=[{}], time index=[{}], filter NaN: [{}]",
369                    self.offset, self.time_index_column_name, self.filter_stale_markers
370                )
371            }
372        }
373    }
374}
375
376pub struct SeriesNormalizeStream {
377    offset: Millisecond,
378    // Column index of TIME INDEX column's position in schema
379    time_index: usize,
380    filter_stale_markers: bool,
381
382    schema: SchemaRef,
383    input: SendableRecordBatchStream,
384    metric: BaselineMetrics,
385    /// Number of series processed.
386    num_series: Count,
387}
388
389impl SeriesNormalizeStream {
390    pub fn normalize(&self, input: RecordBatch) -> DataFusionResult<RecordBatch> {
391        let ts_column = input
392            .column(self.time_index)
393            .as_any()
394            .downcast_ref::<TimestampMillisecondArray>()
395            .ok_or_else(|| {
396                DataFusionError::Execution(
397                    "Time index Column downcast to TimestampMillisecondArray failed".into(),
398                )
399            })?;
400
401        let bias_timestamp = |timestamp: i64| {
402            timestamp.checked_add(self.offset).ok_or_else(|| {
403                DataFusionError::Execution("SeriesNormalize: timestamp offset overflow".into())
404            })
405        };
406
407        // bias the timestamp column by offset
408        let ts_column_biased = if self.offset == 0 {
409            Arc::new(ts_column.clone()) as _
410        } else {
411            Arc::new(ts_column.try_unary::<_, TimestampMillisecondType, _>(&bias_timestamp)?)
412        };
413        let mut columns = input.columns().to_vec();
414        columns[self.time_index] = ts_column_biased;
415
416        // Offset selectors move samples into the evaluation timeline. Keep native histogram
417        // start timestamps on the same timeline for rate and reset calculations.
418        if self.offset != 0 {
419            let native_histogram_type = native_histogram_arrow_type();
420            for column in &mut columns {
421                let Some((histograms, start_timestamp_index, start_timestamps)) = column
422                    .as_any()
423                    .downcast_ref::<StructArray>()
424                    .filter(|histograms| histograms.data_type() == &native_histogram_type)
425                    .and_then(|histograms| {
426                        let (index, _) = histograms.fields().find(START_TIMESTAMP_FIELD)?;
427                        let timestamps = histograms
428                            .column(index)
429                            .as_any()
430                            .downcast_ref::<TimestampMillisecondArray>()?;
431                        Some((histograms, index, timestamps))
432                    })
433                else {
434                    continue;
435                };
436                let start_timestamps = start_timestamps
437                    .try_unary::<_, TimestampMillisecondType, _>(|timestamp| {
438                        // Prometheus uses zero to mean that the start timestamp is unknown.
439                        if timestamp == 0 {
440                            Ok(0)
441                        } else {
442                            bias_timestamp(timestamp)
443                        }
444                    })?;
445                // Replace only the start timestamp child to preserve the histogram payload and
446                // null bitmap.
447                let mut children = histograms.columns().to_vec();
448                children[start_timestamp_index] = Arc::new(start_timestamps);
449                *column = Arc::new(StructArray::new(
450                    histograms.fields().clone(),
451                    children,
452                    histograms.nulls().cloned(),
453                ));
454            }
455        }
456
457        let result_batch = RecordBatch::try_new(input.schema(), columns)?;
458        if !self.filter_stale_markers {
459            return Ok(result_batch);
460        }
461
462        // Filter out Prometheus stale markers.
463        let mut stale_marker_filter = vec![true; input.num_rows()];
464        for column in result_batch.columns() {
465            let Some(stale_sample_column) = prometheus_stale_sample_column(column.as_ref()) else {
466                continue;
467            };
468            for (i, flag) in stale_marker_filter.iter_mut().enumerate() {
469                if is_prometheus_stale_sample(stale_sample_column, i) {
470                    *flag = false;
471                }
472            }
473        }
474
475        let result =
476            compute::filter_record_batch(&result_batch, &BooleanArray::from(stale_marker_filter))
477                .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))?;
478        Ok(result)
479    }
480}
481
482impl RecordBatchStream for SeriesNormalizeStream {
483    fn schema(&self) -> SchemaRef {
484        self.schema.clone()
485    }
486}
487
488impl Stream for SeriesNormalizeStream {
489    type Item = DataFusionResult<RecordBatch>;
490
491    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
492        let poll = match ready!(self.input.poll_next_unpin(cx)) {
493            Some(Ok(batch)) => {
494                self.num_series.add(1);
495                let timer = std::time::Instant::now();
496                let result = Ok(batch).and_then(|batch| self.normalize(batch));
497                self.metric.elapsed_compute().add_elapsed(timer);
498                Poll::Ready(Some(result))
499            }
500            None => {
501                PROMQL_SERIES_COUNT.observe(self.num_series.value() as f64);
502                Poll::Ready(None)
503            }
504            Some(Err(e)) => Poll::Ready(Some(Err(e))),
505        };
506        self.metric.record_poll(poll)
507    }
508}
509
510#[cfg(test)]
511mod test {
512    use common_query::native_histogram::{build_histogram_array, read_histogram};
513    use common_query::prometheus::PROMETHEUS_STALE_NAN_BITS;
514    use datafusion::arrow::array::Float64Array;
515    use datafusion::arrow::buffer::NullBuffer;
516    use datafusion::arrow::datatypes::{
517        ArrowPrimitiveType, DataType, Field, Schema, TimestampMillisecondType,
518    };
519    use datafusion::common::ToDFSchema;
520    use datafusion::datasource::memory::MemorySourceConfig;
521    use datafusion::datasource::source::DataSourceExec;
522    use datafusion::logical_expr::{EmptyRelation, LogicalPlan};
523    use datafusion::prelude::SessionContext;
524    use datatypes::arrow::array::TimestampMillisecondArray;
525    use datatypes::arrow_array::StringArray;
526
527    use super::*;
528    use crate::extension_plan::test_util::native_histogram;
529
530    const TIME_INDEX_COLUMN: &str = "timestamp";
531
532    fn prepare_test_data() -> DataSourceExec {
533        let schema = Arc::new(Schema::new(vec![
534            Field::new(TIME_INDEX_COLUMN, TimestampMillisecondType::DATA_TYPE, true),
535            Field::new("value", DataType::Float64, true),
536            Field::new("path", DataType::Utf8, true),
537        ]));
538        let timestamp_column = Arc::new(TimestampMillisecondArray::from(vec![
539            60_000, 120_000, 0, 30_000, 90_000,
540        ])) as _;
541        let field_column = Arc::new(Float64Array::from(vec![0.0, 1.0, 10.0, 100.0, 1000.0])) as _;
542        let path_column = Arc::new(StringArray::from(vec!["foo", "foo", "foo", "foo", "foo"])) as _;
543        let data = RecordBatch::try_new(
544            schema.clone(),
545            vec![timestamp_column, field_column, path_column],
546        )
547        .unwrap();
548
549        DataSourceExec::new(Arc::new(
550            MemorySourceConfig::try_new(&[vec![data]], schema, None).unwrap(),
551        ))
552    }
553
554    #[test]
555    fn pruning_should_keep_time_and_tag_columns_for_exec() {
556        let df_schema = prepare_test_data().schema().to_dfschema_ref().unwrap();
557        let input = LogicalPlan::EmptyRelation(EmptyRelation {
558            produce_one_row: false,
559            schema: df_schema,
560        });
561        let plan =
562            SeriesNormalize::new(0, TIME_INDEX_COLUMN, true, vec!["path".to_string()], input);
563
564        // Simulate a parent projection requesting only the `value` column.
565        let output_columns = [1usize];
566        let required = plan.necessary_children_exprs(&output_columns).unwrap();
567        let required = &required[0];
568        assert_eq!(required.as_slice(), &[0, 1, 2]);
569    }
570
571    #[tokio::test]
572    async fn test_sort_record_batch() {
573        let memory_exec = Arc::new(prepare_test_data());
574        let normalize_exec = Arc::new(SeriesNormalizeExec {
575            offset: 0,
576            time_index_column_name: TIME_INDEX_COLUMN.to_string(),
577            filter_stale_markers: true,
578            input: memory_exec,
579            tag_columns: vec!["path".to_string()],
580            metric: ExecutionPlanMetricsSet::new(),
581        });
582        let session_context = SessionContext::default();
583        let result = datafusion::physical_plan::collect(normalize_exec, session_context.task_ctx())
584            .await
585            .unwrap();
586        let result_literal = datatypes::arrow::util::pretty::pretty_format_batches(&result)
587            .unwrap()
588            .to_string();
589
590        let expected = String::from(
591            "+---------------------+--------+------+\
592            \n| timestamp           | value  | path |\
593            \n+---------------------+--------+------+\
594            \n| 1970-01-01T00:01:00 | 0.0    | foo  |\
595            \n| 1970-01-01T00:02:00 | 1.0    | foo  |\
596            \n| 1970-01-01T00:00:00 | 10.0   | foo  |\
597            \n| 1970-01-01T00:00:30 | 100.0  | foo  |\
598            \n| 1970-01-01T00:01:30 | 1000.0 | foo  |\
599            \n+---------------------+--------+------+",
600        );
601
602        assert_eq!(result_literal, expected);
603    }
604
605    #[tokio::test]
606    async fn test_offset_record_batch() {
607        let memory_exec = Arc::new(prepare_test_data());
608        let normalize_exec = Arc::new(SeriesNormalizeExec {
609            offset: 1_000,
610            time_index_column_name: TIME_INDEX_COLUMN.to_string(),
611            filter_stale_markers: true,
612            input: memory_exec,
613            metric: ExecutionPlanMetricsSet::new(),
614            tag_columns: vec!["path".to_string()],
615        });
616        let session_context = SessionContext::default();
617        let result = datafusion::physical_plan::collect(normalize_exec, session_context.task_ctx())
618            .await
619            .unwrap();
620        let result_literal = datatypes::arrow::util::pretty::pretty_format_batches(&result)
621            .unwrap()
622            .to_string();
623
624        let expected = String::from(
625            "+---------------------+--------+------+\
626            \n| timestamp           | value  | path |\
627            \n+---------------------+--------+------+\
628            \n| 1970-01-01T00:01:01 | 0.0    | foo  |\
629            \n| 1970-01-01T00:02:01 | 1.0    | foo  |\
630            \n| 1970-01-01T00:00:01 | 10.0   | foo  |\
631            \n| 1970-01-01T00:00:31 | 100.0  | foo  |\
632            \n| 1970-01-01T00:01:31 | 1000.0 | foo  |\
633            \n+---------------------+--------+------+",
634        );
635
636        assert_eq!(result_literal, expected);
637    }
638
639    #[tokio::test]
640    async fn filters_stale_markers_and_preserves_ordinary_nan() {
641        let schema = Arc::new(Schema::new(vec![
642            Field::new(
643                TIME_INDEX_COLUMN,
644                TimestampMillisecondType::DATA_TYPE,
645                false,
646            ),
647            Field::new("value", DataType::Float64, true),
648            Field::new("auxiliary", DataType::Float64, true),
649        ]));
650        let value_column = Float64Array::new(
651            vec![
652                42.0,
653                f64::from_bits(0x7ff0_0000_0000_0002),
654                24.0,
655                f64::from_bits(0x7ff0_0000_0000_0002),
656            ]
657            .into(),
658            Some(NullBuffer::from(vec![true, true, true, false])),
659        );
660        assert!(!value_column.is_valid(3));
661        assert_eq!(value_column.value(3).to_bits(), 0x7ff0_0000_0000_0002);
662        let batch = RecordBatch::try_new(
663            schema.clone(),
664            vec![
665                Arc::new(TimestampMillisecondArray::from(vec![
666                    1_000, 2_000, 3_000, 4_000,
667                ])),
668                Arc::new(value_column),
669                Arc::new(Float64Array::from(vec![
670                    f64::from_bits(0x7ff8_0000_0000_0000),
671                    1.0,
672                    f64::from_bits(0x7ff0_0000_0000_0002),
673                    2.0,
674                ])),
675            ],
676        )
677        .unwrap();
678        let input = Arc::new(DataSourceExec::new(Arc::new(
679            MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(),
680        )));
681        let exec = Arc::new(SeriesNormalizeExec {
682            offset: 0,
683            time_index_column_name: TIME_INDEX_COLUMN.to_string(),
684            filter_stale_markers: true,
685            tag_columns: Vec::new(),
686            input,
687            metric: ExecutionPlanMetricsSet::new(),
688        });
689
690        let context = SessionContext::default();
691        let batches = datafusion::physical_plan::collect(exec, context.task_ctx())
692            .await
693            .unwrap();
694        assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 2);
695        let batch = batches.iter().find(|batch| batch.num_rows() == 2).unwrap();
696        let value = batch
697            .column(1)
698            .as_any()
699            .downcast_ref::<Float64Array>()
700            .unwrap();
701        let auxiliary = batch
702            .column(2)
703            .as_any()
704            .downcast_ref::<Float64Array>()
705            .unwrap();
706
707        assert_eq!(value.value(0), 42.0);
708        assert_eq!(auxiliary.value(0).to_bits(), 0x7ff8_0000_0000_0000);
709        assert!(!value.is_valid(1));
710    }
711
712    #[tokio::test]
713    async fn offsets_native_histogram_timestamps_and_filters_stale_markers() {
714        let mut regular = native_histogram(42.0);
715        regular.start_timestamp = Some(500);
716        let mut ordinary_nan = native_histogram(f64::NAN);
717        ordinary_nan.start_timestamp = Some(0);
718        let histograms = build_histogram_array(&[
719            Some(regular),
720            Some(native_histogram(f64::from_bits(PROMETHEUS_STALE_NAN_BITS))),
721            Some(ordinary_nan),
722            None,
723        ]);
724        let schema = Arc::new(Schema::new(vec![
725            Field::new(
726                TIME_INDEX_COLUMN,
727                TimestampMillisecondType::DATA_TYPE,
728                false,
729            ),
730            Field::new("value", histograms.data_type().clone(), true),
731        ]));
732        let batch = RecordBatch::try_new(
733            schema.clone(),
734            vec![
735                Arc::new(TimestampMillisecondArray::from(vec![
736                    1_000, 2_000, 3_000, 4_000,
737                ])),
738                histograms,
739            ],
740        )
741        .unwrap();
742        let input = Arc::new(DataSourceExec::new(Arc::new(
743            MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(),
744        )));
745        let exec = Arc::new(SeriesNormalizeExec {
746            offset: 1_000,
747            time_index_column_name: TIME_INDEX_COLUMN.to_string(),
748            filter_stale_markers: true,
749            tag_columns: Vec::new(),
750            input,
751            metric: ExecutionPlanMetricsSet::new(),
752        });
753
754        let context = SessionContext::default();
755        let batches = datafusion::physical_plan::collect(exec, context.task_ctx())
756            .await
757            .unwrap();
758        let batch = batches.iter().find(|batch| batch.num_rows() == 3).unwrap();
759        let values = batch
760            .column(1)
761            .as_any()
762            .downcast_ref::<datafusion::arrow::array::StructArray>()
763            .unwrap();
764
765        let timestamps = batch
766            .column(0)
767            .as_any()
768            .downcast_ref::<TimestampMillisecondArray>()
769            .unwrap();
770        assert_eq!(timestamps.values(), &[2_000, 4_000, 5_000]);
771        let regular = read_histogram(values, 0).unwrap().unwrap();
772        assert_eq!((regular.sum, regular.start_timestamp), (42.0, Some(1_500)));
773        let ordinary_nan = read_histogram(values, 1).unwrap().unwrap();
774        assert!(ordinary_nan.sum.is_nan());
775        assert_eq!(ordinary_nan.start_timestamp, Some(0));
776        assert!(read_histogram(values, 2).unwrap().is_none());
777    }
778}