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::pin::Pin;
16use std::sync::Arc;
17use std::task::{Context, Poll};
18
19use common_query::native_histogram::{START_TIMESTAMP_FIELD, native_histogram_arrow_type};
20use datafusion::arrow::array::{Array, BooleanArray, StructArray};
21use datafusion::arrow::compute;
22use datafusion::common::tree_node::TreeNodeRecursion;
23use datafusion::common::{Column, 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    ChildStats, DisplayAs, DisplayFormatType, Distribution, ExecutionPlan,
33    InputDistributionRequirements, PhysicalExpr, PlanProperties, RecordBatchStream,
34    SendableRecordBatchStream, StatisticsArgs,
35};
36use datafusion_expr::ident;
37use datatypes::arrow::array::TimestampMillisecondArray;
38use datatypes::arrow::datatypes::{SchemaRef, TimestampMillisecondType};
39use datatypes::arrow::record_batch::RecordBatch;
40use futures::{Stream, StreamExt, ready};
41use greptime_proto::substrait_extension as pb;
42use prost::Message;
43use snafu::ResultExt;
44
45use crate::error::{DeserializeSnafu, Result};
46use crate::extension_plan::{
47    METRIC_NUM_SERIES, Millisecond, is_prometheus_stale_sample, prometheus_stale_sample_column,
48    resolve_column_name, serialize_column_index,
49};
50use crate::metrics::PROMQL_SERIES_COUNT;
51
52/// Normalizes a single-series input batch and optionally removes Prometheus stale markers.
53///
54/// This node remains the serialized carrier of the selector offset. Native sample timestamps
55/// stay raw: applying an offset can overflow native `i64` ticks even when evaluation is valid.
56/// Manipulators instead apply it in `i128` during selection and produce millisecond outputs.
57/// Histogram start timestamps are already millisecond payloads used for rate/reset, so their
58/// offsets are applied here, preserving unknown zero values and nulls.
59#[derive(Debug, PartialEq, Eq, Hash, PartialOrd)]
60pub struct SeriesNormalize {
61    offset: Millisecond,
62    time_index_column_name: String,
63    filter_stale_markers: bool,
64    tag_columns: Vec<String>,
65
66    input: LogicalPlan,
67    unfix: Option<UnfixIndices>,
68}
69
70#[derive(Debug, PartialEq, Eq, Hash, PartialOrd)]
71struct UnfixIndices {
72    pub time_index_idx: u64,
73    pub tag_column_indices: Vec<u64>,
74}
75
76impl UserDefinedLogicalNodeCore for SeriesNormalize {
77    fn name(&self) -> &str {
78        Self::name()
79    }
80
81    fn inputs(&self) -> Vec<&LogicalPlan> {
82        vec![&self.input]
83    }
84
85    fn schema(&self) -> &DFSchemaRef {
86        self.input.schema()
87    }
88
89    fn expressions(&self) -> Vec<datafusion::logical_expr::Expr> {
90        if self.unfix.is_some() {
91            return vec![];
92        }
93
94        self.tag_columns
95            .iter()
96            .map(ident)
97            .chain(std::iter::once(ident(&self.time_index_column_name)))
98            .collect()
99    }
100
101    fn necessary_children_exprs(&self, output_columns: &[usize]) -> Option<Vec<Vec<usize>>> {
102        if self.unfix.is_some() {
103            return None;
104        }
105
106        let input_schema = self.input.schema();
107        if output_columns.is_empty() {
108            let indices = (0..input_schema.fields().len()).collect::<Vec<_>>();
109            return Some(vec![indices]);
110        }
111
112        let mut required = Vec::with_capacity(output_columns.len() + 1 + self.tag_columns.len());
113        required.extend_from_slice(output_columns);
114        required.push(input_schema.index_of_column_by_name(None, &self.time_index_column_name)?);
115        for tag in &self.tag_columns {
116            required.push(input_schema.index_of_column_by_name(None, tag)?);
117        }
118
119        required.sort_unstable();
120        required.dedup();
121        Some(vec![required])
122    }
123
124    fn fmt_for_explain(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
125        write!(
126            f,
127            "PromSeriesNormalize: offset=[{}], time index=[{}], filter NaN: [{}]",
128            self.offset, self.time_index_column_name, self.filter_stale_markers
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(DataFusionError::Internal(
139                "SeriesNormalize should have at least one input".to_string(),
140            ));
141        }
142
143        let input: LogicalPlan = inputs.into_iter().next().unwrap();
144        let input_schema = input.schema();
145
146        if let Some(unfix) = &self.unfix {
147            // transform indices to names
148            let time_index_column_name = resolve_column_name(
149                unfix.time_index_idx,
150                input_schema,
151                "SeriesNormalize",
152                "time index",
153            )?;
154
155            let tag_columns = unfix
156                .tag_column_indices
157                .iter()
158                .map(|idx| resolve_column_name(*idx, input_schema, "SeriesNormalize", "tag"))
159                .collect::<DataFusionResult<Vec<String>>>()?;
160
161            Ok(Self {
162                offset: self.offset,
163                time_index_column_name,
164                filter_stale_markers: self.filter_stale_markers,
165                tag_columns,
166                input,
167                unfix: None,
168            })
169        } else {
170            Ok(Self {
171                offset: self.offset,
172                time_index_column_name: self.time_index_column_name.clone(),
173                filter_stale_markers: self.filter_stale_markers,
174                tag_columns: self.tag_columns.clone(),
175                input,
176                unfix: None,
177            })
178        }
179    }
180}
181
182impl SeriesNormalize {
183    pub(crate) fn offset_for_time_index(&self, time_index: &Column) -> Option<Millisecond> {
184        let index = self.input.schema().maybe_index_of_column(time_index)?;
185        let (qualifier, field) = self.input.schema().qualified_field(index);
186        (field.name() == &self.time_index_column_name
187            && time_index == &Column::new(qualifier.cloned(), field.name().clone()))
188            .then_some(self.offset)
189    }
190
191    pub fn new<N: AsRef<str>>(
192        offset: Millisecond,
193        time_index_column_name: N,
194        filter_stale_markers: bool,
195        tag_columns: Vec<String>,
196        input: LogicalPlan,
197    ) -> Self {
198        Self {
199            offset,
200            time_index_column_name: time_index_column_name.as_ref().to_string(),
201            filter_stale_markers,
202            tag_columns,
203            input,
204            unfix: None,
205        }
206    }
207
208    pub const fn name() -> &'static str {
209        "SeriesNormalize"
210    }
211
212    /// Returns whether this plan removes Prometheus stale markers.
213    pub const fn filter_stale_markers(&self) -> bool {
214        self.filter_stale_markers
215    }
216
217    pub fn to_execution_plan(&self, exec_input: Arc<dyn ExecutionPlan>) -> Arc<dyn ExecutionPlan> {
218        Arc::new(SeriesNormalizeExec {
219            offset: self.offset,
220            time_index_column_name: self.time_index_column_name.clone(),
221            filter_stale_markers: self.filter_stale_markers,
222            input: exec_input,
223            tag_columns: self.tag_columns.clone(),
224            metric: ExecutionPlanMetricsSet::new(),
225        })
226    }
227
228    pub fn serialize(&self) -> Vec<u8> {
229        let time_index_idx =
230            serialize_column_index(self.input.schema(), &self.time_index_column_name);
231
232        let tag_column_indices = self
233            .tag_columns
234            .iter()
235            .map(|name| serialize_column_index(self.input.schema(), name))
236            .collect::<Vec<u64>>();
237
238        pb::SeriesNormalize {
239            offset: self.offset,
240            time_index_idx,
241            filter_nan: self.filter_stale_markers,
242            tag_column_indices,
243            ..Default::default()
244        }
245        .encode_to_vec()
246    }
247
248    pub fn deserialize(bytes: &[u8]) -> Result<Self> {
249        let pb_normalize = pb::SeriesNormalize::decode(bytes).context(DeserializeSnafu)?;
250        let placeholder_plan = LogicalPlan::EmptyRelation(EmptyRelation {
251            produce_one_row: false,
252            schema: Arc::new(DFSchema::empty()),
253        });
254
255        let unfix = UnfixIndices {
256            time_index_idx: pb_normalize.time_index_idx,
257            tag_column_indices: pb_normalize.tag_column_indices.clone(),
258        };
259
260        Ok(Self {
261            offset: pb_normalize.offset,
262            time_index_column_name: String::new(),
263            filter_stale_markers: pb_normalize.filter_nan,
264            tag_columns: Vec::new(),
265            input: placeholder_plan,
266            unfix: Some(unfix),
267        })
268    }
269}
270
271#[derive(Debug)]
272pub struct SeriesNormalizeExec {
273    offset: Millisecond,
274    time_index_column_name: String,
275    filter_stale_markers: bool,
276    tag_columns: Vec<String>,
277
278    input: Arc<dyn ExecutionPlan>,
279    metric: ExecutionPlanMetricsSet,
280}
281
282impl ExecutionPlan for SeriesNormalizeExec {
283    fn apply_expressions(
284        &self,
285        _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> datafusion_common::Result<TreeNodeRecursion>,
286    ) -> DataFusionResult<TreeNodeRecursion> {
287        Ok(TreeNodeRecursion::Continue)
288    }
289
290    fn schema(&self) -> SchemaRef {
291        self.input.schema()
292    }
293
294    fn input_distribution_requirements(&self) -> InputDistributionRequirements {
295        if self.tag_columns.is_empty() {
296            return InputDistributionRequirements::new(vec![Distribution::SinglePartition]);
297        }
298
299        let schema = self.input.schema();
300        InputDistributionRequirements::new(vec![Distribution::KeyPartitioned(
301            self.tag_columns
302                .iter()
303                // Safety: the tag column names is verified in the planning phase
304                .map(|tag| Arc::new(ColumnExpr::new_with_schema(tag, &schema).unwrap()) as _)
305                .collect(),
306        )])
307    }
308
309    fn properties(&self) -> &Arc<PlanProperties> {
310        self.input.properties()
311    }
312
313    fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
314        vec![&self.input]
315    }
316
317    fn with_new_children(
318        self: Arc<Self>,
319        children: Vec<Arc<dyn ExecutionPlan>>,
320    ) -> DataFusionResult<Arc<dyn ExecutionPlan>> {
321        assert!(!children.is_empty());
322        Ok(Arc::new(Self {
323            offset: self.offset,
324            time_index_column_name: self.time_index_column_name.clone(),
325            filter_stale_markers: self.filter_stale_markers,
326            input: children[0].clone(),
327            tag_columns: self.tag_columns.clone(),
328            metric: self.metric.clone(),
329        }))
330    }
331
332    fn execute(
333        &self,
334        partition: usize,
335        context: Arc<TaskContext>,
336    ) -> DataFusionResult<SendableRecordBatchStream> {
337        let baseline_metric = BaselineMetrics::new(&self.metric, partition);
338        let metrics_builder = MetricBuilder::new(&self.metric);
339        let num_series = Count::new();
340        metrics_builder
341            .with_partition(partition)
342            .build(MetricValue::Count {
343                name: METRIC_NUM_SERIES.into(),
344                count: num_series.clone(),
345            });
346
347        let input = self.input.execute(partition, context)?;
348        let schema = input.schema();
349        Ok(Box::pin(SeriesNormalizeStream {
350            offset: self.offset,
351            filter_stale_markers: self.filter_stale_markers,
352            schema,
353            input,
354            metric: baseline_metric,
355            num_series,
356        }))
357    }
358
359    fn metrics(&self) -> Option<MetricsSet> {
360        Some(self.metric.clone_inner())
361    }
362
363    fn child_stats_requests(&self, partition: Option<usize>) -> Vec<ChildStats> {
364        vec![ChildStats::At(partition)]
365    }
366
367    fn statistics_from_inputs(
368        &self,
369        input_stats: &[Arc<Statistics>],
370        _args: &StatisticsArgs,
371    ) -> DataFusionResult<Arc<Statistics>> {
372        Ok(Arc::clone(&input_stats[0]))
373    }
374
375    fn name(&self) -> &str {
376        "SeriesNormalizeExec"
377    }
378}
379
380impl DisplayAs for SeriesNormalizeExec {
381    fn fmt_as(&self, t: DisplayFormatType, f: &mut std::fmt::Formatter) -> std::fmt::Result {
382        match t {
383            DisplayFormatType::Default
384            | DisplayFormatType::Verbose
385            | DisplayFormatType::TreeRender => {
386                write!(
387                    f,
388                    "PromSeriesNormalizeExec: offset=[{}], time index=[{}], filter NaN: [{}]",
389                    self.offset, self.time_index_column_name, self.filter_stale_markers
390                )
391            }
392        }
393    }
394}
395
396pub struct SeriesNormalizeStream {
397    offset: Millisecond,
398    filter_stale_markers: bool,
399
400    schema: SchemaRef,
401    input: SendableRecordBatchStream,
402    metric: BaselineMetrics,
403    /// Number of series processed.
404    num_series: Count,
405}
406
407impl SeriesNormalizeStream {
408    pub fn normalize(&self, input: RecordBatch) -> DataFusionResult<RecordBatch> {
409        // Native sample timestamps remain raw. Manipulators apply the selector offset
410        // in wide nanosecond arithmetic, avoiding overflow in native Arrow storage.
411        let mut columns = input.columns().to_vec();
412
413        // Offset selectors move native histogram start timestamps onto the evaluation
414        // timeline for rate and reset calculations. These payloads are milliseconds.
415        if self.offset != 0 {
416            let native_histogram_type = native_histogram_arrow_type();
417            for column in &mut columns {
418                let Some((histograms, start_timestamp_index, start_timestamps)) = column
419                    .as_any()
420                    .downcast_ref::<StructArray>()
421                    .filter(|histograms| histograms.data_type() == &native_histogram_type)
422                    .and_then(|histograms| {
423                        let (index, _) = histograms.fields().find(START_TIMESTAMP_FIELD)?;
424                        let timestamps = histograms
425                            .column(index)
426                            .as_any()
427                            .downcast_ref::<TimestampMillisecondArray>()?;
428                        Some((histograms, index, timestamps))
429                    })
430                else {
431                    continue;
432                };
433                let start_timestamps = start_timestamps
434                    .try_unary::<_, TimestampMillisecondType, _>(|timestamp| {
435                        // Prometheus uses zero to mean that the start timestamp is unknown.
436                        if timestamp == 0 {
437                            Ok(0)
438                        } else {
439                            timestamp.checked_add(self.offset).ok_or_else(|| {
440                                DataFusionError::Execution(
441                                    "SeriesNormalize: histogram timestamp offset overflow".into(),
442                                )
443                            })
444                        }
445                    })?;
446                // Struct arrays are immutable, so rebuild the physical histogram payload with
447                // only its start-timestamp child replaced. Its logical schema stays unchanged:
448                // the selector offset affects histogram reset/rate metadata, not the native
449                // sample timestamp column consumed by later manipulators.
450                let mut children = histograms.columns().to_vec();
451                children[start_timestamp_index] = Arc::new(start_timestamps);
452                *column = Arc::new(StructArray::new(
453                    histograms.fields().clone(),
454                    children,
455                    histograms.nulls().cloned(),
456                ));
457            }
458        }
459
460        let result_batch = RecordBatch::try_new(input.schema(), columns)?;
461        if !self.filter_stale_markers {
462            return Ok(result_batch);
463        }
464
465        // Filter out Prometheus stale markers.
466        let mut stale_marker_filter = vec![true; input.num_rows()];
467        for column in result_batch.columns() {
468            let Some(stale_sample_column) = prometheus_stale_sample_column(column.as_ref()) else {
469                continue;
470            };
471            for (i, flag) in stale_marker_filter.iter_mut().enumerate() {
472                if is_prometheus_stale_sample(stale_sample_column, i) {
473                    *flag = false;
474                }
475            }
476        }
477
478        let result =
479            compute::filter_record_batch(&result_batch, &BooleanArray::from(stale_marker_filter))
480                .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))?;
481        Ok(result)
482    }
483}
484
485impl RecordBatchStream for SeriesNormalizeStream {
486    fn schema(&self) -> SchemaRef {
487        self.schema.clone()
488    }
489}
490
491impl Stream for SeriesNormalizeStream {
492    type Item = DataFusionResult<RecordBatch>;
493
494    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
495        let poll = match ready!(self.input.poll_next_unpin(cx)) {
496            Some(Ok(batch)) => {
497                self.num_series.add(1);
498                let timer = std::time::Instant::now();
499                let result = Ok(batch).and_then(|batch| self.normalize(batch));
500                self.metric.elapsed_compute().add_elapsed(timer);
501                Poll::Ready(Some(result))
502            }
503            None => {
504                PROMQL_SERIES_COUNT.observe(self.num_series.value() as f64);
505                Poll::Ready(None)
506            }
507            Some(Err(e)) => Poll::Ready(Some(Err(e))),
508        };
509        self.metric.record_poll(poll)
510    }
511}
512
513#[cfg(test)]
514mod test {
515    use common_query::native_histogram::{build_histogram_array, read_histogram};
516    use common_query::prometheus::PROMETHEUS_STALE_NAN_BITS;
517    use datafusion::arrow::array::{
518        DictionaryArray, Float64Array, TimestampMicrosecondArray, TimestampNanosecondArray,
519    };
520    use datafusion::arrow::buffer::NullBuffer;
521    use datafusion::arrow::datatypes::{
522        ArrowPrimitiveType, DataType, Field, Int64Type, Schema, TimeUnit, TimestampMillisecondType,
523    };
524    use datafusion::common::ToDFSchema;
525    use datafusion::datasource::memory::MemorySourceConfig;
526    use datafusion::datasource::source::DataSourceExec;
527    use datafusion::logical_expr::{EmptyRelation, LogicalPlan};
528    use datafusion::prelude::SessionContext;
529    use datatypes::arrow::array::TimestampMillisecondArray;
530    use datatypes::arrow_array::StringArray;
531
532    use super::*;
533    use crate::extension_plan::RangeManipulate;
534    use crate::extension_plan::test_util::native_histogram;
535    use crate::range_array::RangeArray;
536
537    const TIME_INDEX_COLUMN: &str = "timestamp";
538
539    fn prepare_test_data() -> DataSourceExec {
540        let schema = Arc::new(Schema::new(vec![
541            Field::new(TIME_INDEX_COLUMN, TimestampMillisecondType::DATA_TYPE, true),
542            Field::new("value", DataType::Float64, true),
543            Field::new("path", DataType::Utf8, true),
544        ]));
545        let timestamp_column = Arc::new(TimestampMillisecondArray::from(vec![
546            60_000, 120_000, 0, 30_000, 90_000,
547        ])) as _;
548        let field_column = Arc::new(Float64Array::from(vec![0.0, 1.0, 10.0, 100.0, 1000.0])) as _;
549        let path_column = Arc::new(StringArray::from(vec!["foo", "foo", "foo", "foo", "foo"])) as _;
550        let data = RecordBatch::try_new(
551            schema.clone(),
552            vec![timestamp_column, field_column, path_column],
553        )
554        .unwrap();
555
556        DataSourceExec::new(Arc::new(
557            MemorySourceConfig::try_new(&[vec![data]], schema, None).unwrap(),
558        ))
559    }
560
561    #[test]
562    fn pruning_should_keep_time_and_tag_columns_for_exec() {
563        let df_schema = prepare_test_data().schema().to_dfschema_ref().unwrap();
564        let input = LogicalPlan::EmptyRelation(EmptyRelation {
565            produce_one_row: false,
566            schema: df_schema,
567        });
568        let plan =
569            SeriesNormalize::new(0, TIME_INDEX_COLUMN, true, vec!["path".to_string()], input);
570
571        // Simulate a parent projection requesting only the `value` column.
572        let output_columns = [1usize];
573        let required = plan.necessary_children_exprs(&output_columns).unwrap();
574        let required = &required[0];
575        assert_eq!(required.as_slice(), &[0, 1, 2]);
576    }
577
578    #[tokio::test]
579    async fn test_sort_record_batch() {
580        let memory_exec = Arc::new(prepare_test_data());
581        let normalize_exec = Arc::new(SeriesNormalizeExec {
582            offset: 0,
583            time_index_column_name: TIME_INDEX_COLUMN.to_string(),
584            filter_stale_markers: true,
585            input: memory_exec,
586            tag_columns: vec!["path".to_string()],
587            metric: ExecutionPlanMetricsSet::new(),
588        });
589        let session_context = SessionContext::default();
590        let result = datafusion::physical_plan::collect(normalize_exec, session_context.task_ctx())
591            .await
592            .unwrap();
593        let result_literal = datatypes::arrow::util::pretty::pretty_format_batches(&result)
594            .unwrap()
595            .to_string();
596
597        let expected = String::from(
598            "+---------------------+--------+------+\
599            \n| timestamp           | value  | path |\
600            \n+---------------------+--------+------+\
601            \n| 1970-01-01T00:01:00 | 0.0    | foo  |\
602            \n| 1970-01-01T00:02:00 | 1.0    | foo  |\
603            \n| 1970-01-01T00:00:00 | 10.0   | foo  |\
604            \n| 1970-01-01T00:00:30 | 100.0  | foo  |\
605            \n| 1970-01-01T00:01:30 | 1000.0 | foo  |\
606            \n+---------------------+--------+------+",
607        );
608
609        assert_eq!(result_literal, expected);
610    }
611
612    #[tokio::test]
613    async fn test_offset_record_batch() {
614        let memory_exec = Arc::new(prepare_test_data());
615        let normalize_exec = Arc::new(SeriesNormalizeExec {
616            offset: 1_000,
617            time_index_column_name: TIME_INDEX_COLUMN.to_string(),
618            filter_stale_markers: true,
619            input: memory_exec,
620            metric: ExecutionPlanMetricsSet::new(),
621            tag_columns: vec!["path".to_string()],
622        });
623        let session_context = SessionContext::default();
624        let result = datafusion::physical_plan::collect(normalize_exec, session_context.task_ctx())
625            .await
626            .unwrap();
627        let result_literal = datatypes::arrow::util::pretty::pretty_format_batches(&result)
628            .unwrap()
629            .to_string();
630
631        let expected = String::from(
632            "+---------------------+--------+------+\
633            \n| timestamp           | value  | path |\
634            \n+---------------------+--------+------+\
635            \n| 1970-01-01T00:01:00 | 0.0    | foo  |\
636            \n| 1970-01-01T00:02:00 | 1.0    | foo  |\
637            \n| 1970-01-01T00:00:00 | 10.0   | foo  |\
638            \n| 1970-01-01T00:00:30 | 100.0  | foo  |\
639            \n| 1970-01-01T00:01:30 | 1000.0 | foo  |\
640            \n+---------------------+--------+------+",
641        );
642
643        assert_eq!(result_literal, expected);
644    }
645
646    #[tokio::test]
647    async fn filters_stale_markers_and_preserves_ordinary_nan() {
648        let schema = Arc::new(Schema::new(vec![
649            Field::new(
650                TIME_INDEX_COLUMN,
651                TimestampMillisecondType::DATA_TYPE,
652                false,
653            ),
654            Field::new("value", DataType::Float64, true),
655            Field::new("auxiliary", DataType::Float64, true),
656        ]));
657        let value_column = Float64Array::new(
658            vec![
659                42.0,
660                f64::from_bits(0x7ff0_0000_0000_0002),
661                24.0,
662                f64::from_bits(0x7ff0_0000_0000_0002),
663            ]
664            .into(),
665            Some(NullBuffer::from(vec![true, true, true, false])),
666        );
667        assert!(!value_column.is_valid(3));
668        assert_eq!(value_column.value(3).to_bits(), 0x7ff0_0000_0000_0002);
669        let batch = RecordBatch::try_new(
670            schema.clone(),
671            vec![
672                Arc::new(TimestampMillisecondArray::from(vec![
673                    1_000, 2_000, 3_000, 4_000,
674                ])),
675                Arc::new(value_column),
676                Arc::new(Float64Array::from(vec![
677                    f64::from_bits(0x7ff8_0000_0000_0000),
678                    1.0,
679                    f64::from_bits(0x7ff0_0000_0000_0002),
680                    2.0,
681                ])),
682            ],
683        )
684        .unwrap();
685        let input = Arc::new(DataSourceExec::new(Arc::new(
686            MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(),
687        )));
688        let exec = Arc::new(SeriesNormalizeExec {
689            offset: 0,
690            time_index_column_name: TIME_INDEX_COLUMN.to_string(),
691            filter_stale_markers: true,
692            tag_columns: Vec::new(),
693            input,
694            metric: ExecutionPlanMetricsSet::new(),
695        });
696
697        let context = SessionContext::default();
698        let batches = datafusion::physical_plan::collect(exec, context.task_ctx())
699            .await
700            .unwrap();
701        assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 2);
702        let batch = batches.iter().find(|batch| batch.num_rows() == 2).unwrap();
703        let value = batch
704            .column(1)
705            .as_any()
706            .downcast_ref::<Float64Array>()
707            .unwrap();
708        let auxiliary = batch
709            .column(2)
710            .as_any()
711            .downcast_ref::<Float64Array>()
712            .unwrap();
713
714        assert_eq!(value.value(0), 42.0);
715        assert_eq!(auxiliary.value(0).to_bits(), 0x7ff8_0000_0000_0000);
716        assert!(!value.is_valid(1));
717    }
718
719    #[tokio::test]
720    async fn offsets_native_histogram_timestamps_and_filters_stale_markers() {
721        let mut regular = native_histogram(42.0);
722        regular.start_timestamp = Some(500);
723        let mut ordinary_nan = native_histogram(f64::NAN);
724        ordinary_nan.start_timestamp = Some(0);
725        let mut unknown_start = native_histogram(7.0);
726        unknown_start.start_timestamp = None;
727        let histograms = build_histogram_array(&[
728            Some(regular),
729            Some(native_histogram(f64::from_bits(PROMETHEUS_STALE_NAN_BITS))),
730            Some(ordinary_nan),
731            Some(unknown_start),
732            None,
733        ]);
734        for (unit, ticks_per_ms) in [
735            (TimeUnit::Millisecond, 1_i64),
736            (TimeUnit::Microsecond, 1_000),
737            (TimeUnit::Nanosecond, 1_000_000),
738        ] {
739            let timestamp_array = |values: Vec<i64>| -> Arc<dyn Array> {
740                match unit {
741                    TimeUnit::Millisecond => Arc::new(TimestampMillisecondArray::from(values)),
742                    TimeUnit::Microsecond => Arc::new(TimestampMicrosecondArray::from(values)),
743                    TimeUnit::Nanosecond => Arc::new(TimestampNanosecondArray::from(values)),
744                    TimeUnit::Second => unreachable!(),
745                }
746            };
747            for offset in [-1_i64, 1] {
748                let timestamps = timestamp_array(
749                    [1_000, 2_000, 3_000, 4_000, 5_000]
750                        .into_iter()
751                        .map(|timestamp| timestamp * ticks_per_ms)
752                        .collect(),
753                );
754                let schema = Arc::new(Schema::new(vec![
755                    Field::new(TIME_INDEX_COLUMN, timestamps.data_type().clone(), false),
756                    Field::new("value", histograms.data_type().clone(), true),
757                ]));
758                let batch =
759                    RecordBatch::try_new(schema.clone(), vec![timestamps, histograms.clone()])
760                        .unwrap();
761                let input = Arc::new(DataSourceExec::new(Arc::new(
762                    MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(),
763                )));
764                let exec = Arc::new(SeriesNormalizeExec {
765                    offset,
766                    time_index_column_name: TIME_INDEX_COLUMN.to_string(),
767                    filter_stale_markers: true,
768                    tag_columns: Vec::new(),
769                    input,
770                    metric: ExecutionPlanMetricsSet::new(),
771                });
772                let context = SessionContext::default();
773                let batches = datafusion::physical_plan::collect(exec, context.task_ctx())
774                    .await
775                    .unwrap();
776                assert_eq!(
777                    batches.iter().map(RecordBatch::num_rows).sum::<usize>(),
778                    4,
779                    "unit={unit:?}, offset={offset}"
780                );
781                let batch = batches.iter().find(|batch| batch.num_rows() == 4).unwrap();
782                let expected_timestamps = timestamp_array(
783                    [1_000, 3_000, 4_000, 5_000]
784                        .into_iter()
785                        .map(|timestamp| timestamp * ticks_per_ms)
786                        .collect(),
787                );
788                assert_eq!(
789                    batch.column(0).to_data(),
790                    expected_timestamps.to_data(),
791                    "unit={unit:?}, offset={offset}"
792                );
793                let values = batch
794                    .column(1)
795                    .as_any()
796                    .downcast_ref::<datafusion::arrow::array::StructArray>()
797                    .unwrap();
798                let regular = read_histogram(values, 0).unwrap().unwrap();
799                assert_eq!(
800                    (regular.sum, regular.start_timestamp),
801                    (42.0, Some(500 + offset)),
802                    "unit={unit:?}, offset={offset}"
803                );
804                let ordinary_nan = read_histogram(values, 1).unwrap().unwrap();
805                assert!(ordinary_nan.sum.is_nan());
806                assert_eq!(ordinary_nan.start_timestamp, Some(0));
807                let unknown_start = read_histogram(values, 2).unwrap().unwrap();
808                assert_eq!(
809                    (unknown_start.sum, unknown_start.start_timestamp),
810                    (7.0, None)
811                );
812                assert!(read_histogram(values, 3).unwrap().is_none());
813            }
814        }
815
816        let mut known_start = native_histogram(42.0);
817        known_start.start_timestamp = Some(500);
818        let mut sentinel_start = native_histogram(8.0);
819        sentinel_start.start_timestamp = Some(0);
820        let unknown_start = native_histogram(7.0);
821        let histograms =
822            build_histogram_array(&[Some(known_start), Some(sentinel_start), Some(unknown_start)]);
823        let schema = Arc::new(Schema::new(vec![
824            Field::new(
825                TIME_INDEX_COLUMN,
826                TimestampMillisecondType::DATA_TYPE,
827                false,
828            ),
829            Field::new("value", histograms.data_type().clone(), true),
830        ]));
831        let batch = RecordBatch::try_new(
832            schema.clone(),
833            vec![
834                Arc::new(TimestampMillisecondArray::from(vec![1_000; 3])),
835                histograms,
836            ],
837        )
838        .unwrap();
839        let logical_input = LogicalPlan::EmptyRelation(EmptyRelation {
840            produce_one_row: false,
841            schema: schema.clone().to_dfschema_ref().unwrap(),
842        });
843        let normalized =
844            SeriesNormalize::new(1_000, TIME_INDEX_COLUMN, false, Vec::new(), logical_input);
845        let range = RangeManipulate::new(
846            2_000,
847            2_000,
848            1,
849            1_000,
850            1,
851            TIME_INDEX_COLUMN.to_string(),
852            vec!["value".to_string()],
853            LogicalPlan::Extension(datafusion::logical_expr::Extension {
854                node: Arc::new(normalized),
855            }),
856        )
857        .unwrap();
858        let input = Arc::new(DataSourceExec::new(Arc::new(
859            MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(),
860        )));
861        let normalized_input = Arc::new(SeriesNormalizeExec {
862            offset: 1_000,
863            time_index_column_name: TIME_INDEX_COLUMN.to_string(),
864            filter_stale_markers: false,
865            tag_columns: Vec::new(),
866            input,
867            metric: ExecutionPlanMetricsSet::new(),
868        });
869        let output = datafusion::physical_plan::collect(
870            range.to_execution_plan(normalized_input),
871            SessionContext::default().task_ctx(),
872        )
873        .await
874        .unwrap();
875        let values = RangeArray::try_new(
876            output[0]
877                .column(1)
878                .as_any()
879                .downcast_ref::<DictionaryArray<Int64Type>>()
880                .unwrap()
881                .clone(),
882        )
883        .unwrap();
884        let values = values.get(0).unwrap();
885        let values = values
886            .as_any()
887            .downcast_ref::<datafusion::arrow::array::StructArray>()
888            .unwrap();
889        assert_eq!(
890            read_histogram(values, 0).unwrap().unwrap().start_timestamp,
891            Some(1_500)
892        );
893        assert_eq!(
894            read_histogram(values, 1).unwrap().unwrap().start_timestamp,
895            Some(0)
896        );
897        assert_eq!(
898            read_histogram(values, 2).unwrap().unwrap().start_timestamp,
899            None
900        );
901        let timestamps = RangeArray::try_new(
902            output[0]
903                .column(2)
904                .as_any()
905                .downcast_ref::<DictionaryArray<Int64Type>>()
906                .unwrap()
907                .clone(),
908        )
909        .unwrap();
910        assert_eq!(
911            timestamps.get(0).unwrap().to_data(),
912            TimestampMillisecondArray::from(vec![2_000; 3]).to_data()
913        );
914    }
915}