1use std::collections::HashMap;
16use std::ops::Div;
17use std::pin::Pin;
18use std::sync::Arc;
19use std::task::{Context, Poll};
20
21use datafusion::arrow::array::ArrayRef;
22use datafusion::arrow::datatypes::{DataType, TimeUnit};
23use datafusion::catalog::Session;
24use datafusion::common::arrow::datatypes::Field;
25use datafusion::common::stats::Precision;
26use datafusion::common::tree_node::TreeNodeRecursion;
27use datafusion::common::{
28 DFSchema, DFSchemaRef, Result as DataFusionResult, Statistics, TableReference,
29};
30use datafusion::datasource::{MemTable, provider_as_source};
31use datafusion::error::DataFusionError;
32use datafusion::execution::context::TaskContext;
33use datafusion::logical_expr::physical_planning_context::PhysicalPlanningContext;
34use datafusion::logical_expr::{ExprSchemable, LogicalPlan, UserDefinedLogicalNodeCore};
35use datafusion::physical_expr::{EquivalenceProperties, PhysicalExpr, PhysicalExprRef};
36use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
37use datafusion::physical_plan::metrics::{BaselineMetrics, ExecutionPlanMetricsSet, MetricsSet};
38use datafusion::physical_plan::{
39 DisplayAs, DisplayFormatType, ExecutionPlan, Partitioning, PlanProperties, RecordBatchStream,
40 SendableRecordBatchStream, StatisticsArgs,
41};
42use datafusion::physical_planner::PhysicalPlanner;
43use datafusion::prelude::{Expr, ident, lit};
44use datafusion_expr::LogicalPlanBuilder;
45use datatypes::arrow::array::TimestampMillisecondArray;
46use datatypes::arrow::datatypes::SchemaRef;
47use datatypes::arrow::record_batch::RecordBatch;
48use futures::Stream;
49
50use crate::extension_plan::Millisecond;
51
52#[derive(Debug, Clone, PartialEq, Eq, Hash)]
57pub struct EmptyMetric {
58 start: Millisecond,
59 end: Millisecond,
60 interval: Millisecond,
61 expr: Option<Expr>,
62 time_index_schema: DFSchemaRef,
65 result_schema: DFSchemaRef,
67 dummy_input: LogicalPlan,
73}
74
75impl EmptyMetric {
76 pub fn new(
77 start: Millisecond,
78 end: Millisecond,
79 interval: Millisecond,
80 time_index_column_name: String,
81 field_column_name: String,
82 field_expr: Option<Expr>,
83 ) -> DataFusionResult<Self> {
84 let qualifier = Some(TableReference::bare(""));
85 let ts_only_schema = build_ts_only_schema(&time_index_column_name);
86 let mut fields = vec![(qualifier.clone(), ts_only_schema.field(0).clone())];
87 if let Some(field_expr) = &field_expr {
88 let field_data_type = field_expr.get_type(&ts_only_schema)?;
89 fields.push((
90 qualifier.clone(),
91 Arc::new(Field::new(field_column_name, field_data_type, true)),
92 ));
93 }
94 let schema = Arc::new(DFSchema::new_with_metadata(fields, HashMap::new())?);
95
96 let table = MemTable::try_new(Arc::new(schema.as_arrow().clone()), vec![vec![]])?;
97 let source = provider_as_source(Arc::new(table));
98 let dummy_input =
99 LogicalPlanBuilder::scan("dummy", source, None).and_then(|x| x.build())?;
100
101 Ok(Self {
102 start,
103 end,
104 interval,
105 time_index_schema: Arc::new(ts_only_schema),
106 result_schema: schema,
107 expr: field_expr,
108 dummy_input,
109 })
110 }
111
112 pub const fn name() -> &'static str {
113 "EmptyMetric"
114 }
115
116 pub fn to_execution_plan(
117 &self,
118 session: &dyn Session,
119 physical_planner: &dyn PhysicalPlanner,
120 planning_ctx: &PhysicalPlanningContext,
121 ) -> DataFusionResult<Arc<dyn ExecutionPlan>> {
122 let physical_expr = self
123 .expr
124 .as_ref()
125 .map(|expr| {
126 physical_planner.create_physical_expr(
127 expr,
128 &self.time_index_schema,
129 session,
130 planning_ctx,
131 )
132 })
133 .transpose()?;
134 let result_schema: SchemaRef = self.result_schema.inner().clone();
135 let properties = Arc::new(PlanProperties::new(
136 EquivalenceProperties::new(result_schema.clone()),
137 Partitioning::UnknownPartitioning(1),
138 EmissionType::Incremental,
139 Boundedness::Bounded,
140 ));
141 Ok(Arc::new(EmptyMetricExec {
142 start: self.start,
143 end: self.end,
144 interval: self.interval,
145 time_index_schema: self.time_index_schema.inner().clone(),
146 result_schema,
147 expr: physical_expr,
148 properties,
149 metric: ExecutionPlanMetricsSet::new(),
150 }))
151 }
152}
153
154impl UserDefinedLogicalNodeCore for EmptyMetric {
155 fn name(&self) -> &str {
156 Self::name()
157 }
158
159 fn inputs(&self) -> Vec<&LogicalPlan> {
160 vec![&self.dummy_input]
161 }
162
163 fn schema(&self) -> &DFSchemaRef {
164 &self.result_schema
165 }
166
167 fn expressions(&self) -> Vec<Expr> {
168 if let Some(expr) = &self.expr {
169 vec![expr.clone()]
170 } else {
171 vec![]
172 }
173 }
174
175 fn fmt_for_explain(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
176 write!(
177 f,
178 "EmptyMetric: range=[{}..{}], interval=[{}]",
179 self.start, self.end, self.interval,
180 )
181 }
182
183 fn with_exprs_and_inputs(
184 &self,
185 exprs: Vec<Expr>,
186 _inputs: Vec<LogicalPlan>,
187 ) -> DataFusionResult<Self> {
188 Ok(Self {
189 start: self.start,
190 end: self.end,
191 interval: self.interval,
192 expr: exprs.into_iter().next(),
193 time_index_schema: self.time_index_schema.clone(),
194 result_schema: self.result_schema.clone(),
195 dummy_input: self.dummy_input.clone(),
196 })
197 }
198}
199
200impl PartialOrd for EmptyMetric {
201 fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
202 match self.start.partial_cmp(&other.start) {
204 Some(core::cmp::Ordering::Equal) => {}
205 ord => return ord,
206 }
207 match self.end.partial_cmp(&other.end) {
208 Some(core::cmp::Ordering::Equal) => {}
209 ord => return ord,
210 }
211 match self.interval.partial_cmp(&other.interval) {
212 Some(core::cmp::Ordering::Equal) => {}
213 ord => return ord,
214 }
215 self.expr.partial_cmp(&other.expr)
216 }
217}
218
219#[derive(Debug, Clone)]
220pub struct EmptyMetricExec {
221 start: Millisecond,
222 end: Millisecond,
223 interval: Millisecond,
224 time_index_schema: SchemaRef,
227 result_schema: SchemaRef,
229 expr: Option<PhysicalExprRef>,
230 properties: Arc<PlanProperties>,
231 metric: ExecutionPlanMetricsSet,
232}
233
234impl ExecutionPlan for EmptyMetricExec {
235 fn apply_expressions(
236 &self,
237 f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> datafusion_common::Result<TreeNodeRecursion>,
238 ) -> DataFusionResult<TreeNodeRecursion> {
239 datafusion::physical_plan::apply_expression_roots(self.expr.iter(), f)
240 }
241
242 fn schema(&self) -> SchemaRef {
243 self.result_schema.clone()
244 }
245
246 fn properties(&self) -> &Arc<PlanProperties> {
247 &self.properties
248 }
249
250 fn maintains_input_order(&self) -> Vec<bool> {
251 vec![]
252 }
253
254 fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
255 vec![]
256 }
257
258 fn with_new_children(
259 self: Arc<Self>,
260 _children: Vec<Arc<dyn ExecutionPlan>>,
261 ) -> DataFusionResult<Arc<dyn ExecutionPlan>> {
262 Ok(Arc::new(self.as_ref().clone()))
263 }
264
265 fn execute(
266 &self,
267 partition: usize,
268 _context: Arc<TaskContext>,
269 ) -> DataFusionResult<SendableRecordBatchStream> {
270 let baseline_metric = BaselineMetrics::new(&self.metric, partition);
271 Ok(Box::pin(EmptyMetricStream {
272 start: self.start,
273 end: self.end,
274 interval: self.interval,
275 expr: self.expr.clone(),
276 is_first_poll: true,
277 time_index_schema: self.time_index_schema.clone(),
278 result_schema: self.result_schema.clone(),
279 metric: baseline_metric,
280 }))
281 }
282
283 fn metrics(&self) -> Option<MetricsSet> {
284 Some(self.metric.clone_inner())
285 }
286
287 fn statistics_from_inputs(
288 &self,
289 _input_stats: &[Arc<Statistics>],
290 args: &StatisticsArgs,
291 ) -> DataFusionResult<Arc<Statistics>> {
292 let partition = args.partition();
293 if partition.is_some() {
294 return Ok(Arc::new(Statistics::new_unknown(self.schema().as_ref())));
295 }
296
297 let estimated_row_num = if self.end > self.start {
298 (self.end - self.start) as f64 / self.interval as f64
299 } else {
300 0.0
301 };
302 let total_byte_size = estimated_row_num * std::mem::size_of::<Millisecond>() as f64;
303
304 Ok(Arc::new(Statistics {
305 num_rows: Precision::Inexact(estimated_row_num.floor() as _),
306 total_byte_size: Precision::Inexact(total_byte_size.floor() as _),
307 column_statistics: Statistics::unknown_column(&self.schema()),
308 }))
309 }
310
311 fn name(&self) -> &str {
312 "EmptyMetricExec"
313 }
314}
315
316impl DisplayAs for EmptyMetricExec {
317 fn fmt_as(&self, t: DisplayFormatType, f: &mut std::fmt::Formatter) -> std::fmt::Result {
318 match t {
319 DisplayFormatType::Default
320 | DisplayFormatType::Verbose
321 | DisplayFormatType::TreeRender => write!(
322 f,
323 "EmptyMetric: range=[{}..{}], interval=[{}]",
324 self.start, self.end, self.interval,
325 ),
326 }
327 }
328}
329
330pub struct EmptyMetricStream {
331 start: Millisecond,
332 end: Millisecond,
333 interval: Millisecond,
334 expr: Option<PhysicalExprRef>,
335 is_first_poll: bool,
337 time_index_schema: SchemaRef,
340 result_schema: SchemaRef,
342 metric: BaselineMetrics,
343}
344
345impl RecordBatchStream for EmptyMetricStream {
346 fn schema(&self) -> SchemaRef {
347 self.result_schema.clone()
348 }
349}
350
351impl Stream for EmptyMetricStream {
352 type Item = DataFusionResult<RecordBatch>;
353
354 fn poll_next(mut self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
355 let result = if self.is_first_poll {
356 self.is_first_poll = false;
357 let _timer = self.metric.elapsed_compute().timer();
358
359 let time_array = (self.start..=self.end)
362 .step_by(self.interval as _)
363 .collect::<Vec<_>>();
364 let time_array = Arc::new(TimestampMillisecondArray::from(time_array));
365 let num_rows = time_array.len();
366 let input_record_batch =
367 RecordBatch::try_new(self.time_index_schema.clone(), vec![time_array.clone()])
368 .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))?;
369 let mut result_arrays: Vec<ArrayRef> = vec![time_array];
370
371 if let Some(field_expr) = &self.expr {
373 result_arrays.push(
374 field_expr
375 .evaluate(&input_record_batch)
376 .and_then(|x| x.into_array(num_rows))?,
377 );
378 }
379
380 let batch = RecordBatch::try_new(self.result_schema.clone(), result_arrays)
382 .map_err(|e| DataFusionError::ArrowError(Box::new(e), None));
383
384 Poll::Ready(Some(batch))
385 } else {
386 Poll::Ready(None)
387 };
388 self.metric.record_poll(result)
389 }
390}
391
392fn build_ts_only_schema(column_name: &str) -> DFSchema {
394 let ts_field = Field::new(
395 column_name,
396 DataType::Timestamp(TimeUnit::Millisecond, None),
397 false,
398 );
399 DFSchema::new_with_metadata(
401 vec![(Some(TableReference::bare("")), Arc::new(ts_field))],
402 HashMap::new(),
403 )
404 .unwrap()
405}
406
407pub fn build_special_time_expr(time_index_column_name: &str) -> Expr {
410 let input_schema = build_ts_only_schema(time_index_column_name);
411 ident(time_index_column_name)
413 .cast_to(&DataType::Int64, &input_schema)
414 .unwrap()
415 .cast_to(&DataType::Float64, &input_schema)
416 .unwrap()
417 .div(lit(1000.0)) }
419
420#[cfg(test)]
421mod test {
422 use datafusion::physical_planner::DefaultPhysicalPlanner;
423 use datafusion::prelude::SessionContext;
424
425 use super::*;
426
427 async fn do_empty_metric_test(
428 start: Millisecond,
429 end: Millisecond,
430 interval: Millisecond,
431 time_column_name: String,
432 field_column_name: String,
433 expected: String,
434 ) {
435 let session_context = SessionContext::default();
436 let df_default_physical_planner = DefaultPhysicalPlanner::default();
437 let time_expr = build_special_time_expr(&time_column_name);
438 let empty_metric = EmptyMetric::new(
439 start,
440 end,
441 interval,
442 time_column_name,
443 field_column_name,
444 Some(time_expr),
445 )
446 .unwrap();
447 let empty_metric_exec = empty_metric
448 .to_execution_plan(
449 &session_context.state(),
450 &df_default_physical_planner,
451 &PhysicalPlanningContext::default(),
452 )
453 .unwrap();
454
455 let result =
456 datafusion::physical_plan::collect(empty_metric_exec, session_context.task_ctx())
457 .await
458 .unwrap();
459 let result_literal = datatypes::arrow::util::pretty::pretty_format_batches(&result)
460 .unwrap()
461 .to_string();
462
463 assert_eq!(result_literal, expected);
464 }
465
466 #[tokio::test]
467 async fn normal_empty_metric_test() {
468 do_empty_metric_test(
469 0,
470 100,
471 10,
472 "time".to_string(),
473 "value".to_string(),
474 String::from(
475 "+-------------------------+-------+\
476 \n| time | value |\
477 \n+-------------------------+-------+\
478 \n| 1970-01-01T00:00:00 | 0.0 |\
479 \n| 1970-01-01T00:00:00.010 | 0.01 |\
480 \n| 1970-01-01T00:00:00.020 | 0.02 |\
481 \n| 1970-01-01T00:00:00.030 | 0.03 |\
482 \n| 1970-01-01T00:00:00.040 | 0.04 |\
483 \n| 1970-01-01T00:00:00.050 | 0.05 |\
484 \n| 1970-01-01T00:00:00.060 | 0.06 |\
485 \n| 1970-01-01T00:00:00.070 | 0.07 |\
486 \n| 1970-01-01T00:00:00.080 | 0.08 |\
487 \n| 1970-01-01T00:00:00.090 | 0.09 |\
488 \n| 1970-01-01T00:00:00.100 | 0.1 |\
489 \n+-------------------------+-------+",
490 ),
491 )
492 .await
493 }
494
495 #[tokio::test]
496 async fn unaligned_empty_metric_test() {
497 do_empty_metric_test(
498 0,
499 100,
500 11,
501 "time".to_string(),
502 "value".to_string(),
503 String::from(
504 "+-------------------------+-------+\
505 \n| time | value |\
506 \n+-------------------------+-------+\
507 \n| 1970-01-01T00:00:00 | 0.0 |\
508 \n| 1970-01-01T00:00:00.011 | 0.011 |\
509 \n| 1970-01-01T00:00:00.022 | 0.022 |\
510 \n| 1970-01-01T00:00:00.033 | 0.033 |\
511 \n| 1970-01-01T00:00:00.044 | 0.044 |\
512 \n| 1970-01-01T00:00:00.055 | 0.055 |\
513 \n| 1970-01-01T00:00:00.066 | 0.066 |\
514 \n| 1970-01-01T00:00:00.077 | 0.077 |\
515 \n| 1970-01-01T00:00:00.088 | 0.088 |\
516 \n| 1970-01-01T00:00:00.099 | 0.099 |\
517 \n+-------------------------+-------+",
518 ),
519 )
520 .await
521 }
522
523 #[tokio::test]
524 async fn one_row_empty_metric_test() {
525 do_empty_metric_test(
526 0,
527 100,
528 1000,
529 "time".to_string(),
530 "value".to_string(),
531 String::from(
532 "+---------------------+-------+\
533 \n| time | value |\
534 \n+---------------------+-------+\
535 \n| 1970-01-01T00:00:00 | 0.0 |\
536 \n+---------------------+-------+",
537 ),
538 )
539 .await
540 }
541
542 #[tokio::test]
543 async fn negative_range_empty_metric_test() {
544 do_empty_metric_test(
545 1000,
546 -1000,
547 10,
548 "time".to_string(),
549 "value".to_string(),
550 String::from(
551 "+------+-------+\
552 \n| time | value |\
553 \n+------+-------+\
554 \n+------+-------+",
555 ),
556 )
557 .await
558 }
559
560 #[tokio::test]
561 async fn no_field_expr() {
562 let session_context = SessionContext::default();
563 let df_default_physical_planner = DefaultPhysicalPlanner::default();
564 let empty_metric =
565 EmptyMetric::new(0, 200, 1000, "time".to_string(), "value".to_string(), None).unwrap();
566 let empty_metric_exec = empty_metric
567 .to_execution_plan(
568 &session_context.state(),
569 &df_default_physical_planner,
570 &PhysicalPlanningContext::default(),
571 )
572 .unwrap();
573
574 let result =
575 datafusion::physical_plan::collect(empty_metric_exec, session_context.task_ctx())
576 .await
577 .unwrap();
578 let result_literal = datatypes::arrow::util::pretty::pretty_format_batches(&result)
579 .unwrap()
580 .to_string();
581
582 let expected = String::from(
583 "+---------------------+\
584 \n| time |\
585 \n+---------------------+\
586 \n| 1970-01-01T00:00:00 |\
587 \n+---------------------+",
588 );
589 assert_eq!(result_literal, expected);
590 }
591}