1use 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#[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 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 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 .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 num_series: Count,
405}
406
407impl SeriesNormalizeStream {
408 pub fn normalize(&self, input: RecordBatch) -> DataFusionResult<RecordBatch> {
409 let mut columns = input.columns().to_vec();
412
413 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 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 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 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 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}