1use 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#[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 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 .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 time_index: usize,
380 filter_stale_markers: bool,
381
382 schema: SchemaRef,
383 input: SendableRecordBatchStream,
384 metric: BaselineMetrics,
385 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 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 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 if timestamp == 0 {
440 Ok(0)
441 } else {
442 bias_timestamp(timestamp)
443 }
444 })?;
445 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 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 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}