1use std::fmt::{self, Display};
16use std::future::Future;
17use std::marker::PhantomData;
18use std::pin::Pin;
19use std::str::FromStr;
20use std::sync::Arc;
21use std::sync::atomic::{AtomicU64, Ordering};
22use std::task::{Context, Poll};
23
24use common_base::readable_size::ReadableSize;
25use common_telemetry::tracing::{Span, info_span};
26use common_time::util::format_nanoseconds_human_readable;
27use datafusion::arrow::compute::cast;
28use datafusion::arrow::datatypes::SchemaRef as DfSchemaRef;
29use datafusion::error::Result as DfResult;
30use datafusion::execution::context::ExecutionProps;
31use datafusion::logical_expr::Expr;
32use datafusion::logical_expr::utils::conjunction;
33use datafusion::physical_expr::create_physical_expr;
34use datafusion::physical_plan::metrics::{BaselineMetrics, MetricValue};
35use datafusion::physical_plan::{
36 DisplayFormatType, ExecutionPlan, ExecutionPlanVisitor, PhysicalExpr,
37 RecordBatchStream as DfRecordBatchStream, accept,
38};
39use datafusion_common::arrow::error::ArrowError;
40use datafusion_common::{DataFusionError, ToDFSchema};
41use datatypes::arrow::array::Array;
42use datatypes::arrow::datatypes::DataType as ArrowDataType;
43use datatypes::schema::{ColumnExtType, Schema, SchemaRef};
44use futures::ready;
45use jsonb;
46use pin_project::pin_project;
47use snafu::ResultExt;
48
49use crate::error::{self, Result};
50use crate::filter::batch_filter;
51use crate::{
52 DfRecordBatch, DfSendableRecordBatchStream, OrderOption, RecordBatch, RecordBatchStream,
53 SendableRecordBatchStream, Stream,
54};
55
56const REGION_SCAN_EXEC_NAME: &str = "RegionScanExec";
57
58type FutureStream =
59 Pin<Box<dyn std::future::Future<Output = Result<SendableRecordBatchStream>> + Send>>;
60
61#[pin_project]
63pub struct RecordBatchStreamTypeAdapter<T, E> {
64 #[pin]
65 stream: T,
66 projected_schema: DfSchemaRef,
67 projection: Vec<usize>,
68 predicate: Option<Arc<dyn PhysicalExpr>>,
69 phantom: PhantomData<E>,
70}
71
72impl<T, E> RecordBatchStreamTypeAdapter<T, E>
73where
74 T: Stream<Item = std::result::Result<DfRecordBatch, E>>,
75 E: std::error::Error + Send + Sync + 'static,
76{
77 pub fn new(projected_schema: DfSchemaRef, stream: T, projection: Option<Vec<usize>>) -> Self {
78 let projection = if let Some(projection) = projection {
79 projection
80 } else {
81 (0..projected_schema.fields().len()).collect()
82 };
83
84 Self {
85 stream,
86 projected_schema,
87 projection,
88 predicate: None,
89 phantom: Default::default(),
90 }
91 }
92
93 pub fn with_filter(mut self, filters: Vec<Expr>) -> Result<Self> {
94 let filters = if let Some(expr) = conjunction(filters) {
95 let df_schema = self
96 .projected_schema
97 .clone()
98 .to_dfschema_ref()
99 .context(error::PhysicalExprSnafu)?;
100
101 let filters = create_physical_expr(&expr, &df_schema, &ExecutionProps::new())
102 .context(error::PhysicalExprSnafu)?;
103 Some(filters)
104 } else {
105 None
106 };
107 self.predicate = filters;
108 Ok(self)
109 }
110}
111
112impl<T, E> DfRecordBatchStream for RecordBatchStreamTypeAdapter<T, E>
113where
114 T: Stream<Item = std::result::Result<DfRecordBatch, E>>,
115 E: std::error::Error + Send + Sync + 'static,
116{
117 fn schema(&self) -> DfSchemaRef {
118 self.projected_schema.clone()
119 }
120}
121
122impl<T, E> Stream for RecordBatchStreamTypeAdapter<T, E>
123where
124 T: Stream<Item = std::result::Result<DfRecordBatch, E>>,
125 E: std::error::Error + Send + Sync + 'static,
126{
127 type Item = DfResult<DfRecordBatch>;
128
129 fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
130 let this = self.project();
131
132 let batch = futures::ready!(this.stream.poll_next(cx))
133 .map(|r| r.map_err(|e| DataFusionError::External(Box::new(e))));
134
135 let projected_schema = this.projected_schema.clone();
136 let projection = this.projection.clone();
137 let predicate = this.predicate.clone();
138
139 let batch = batch.map(|b| {
140 b.and_then(|b| {
141 let projected_column = b.project(&projection)?;
142 if projected_column.schema().fields.len() != projected_schema.fields.len() {
143 return Err(DataFusionError::ArrowError(Box::new(ArrowError::SchemaError(format!(
144 "Trying to cast a RecordBatch into an incompatible schema. RecordBatch: {}, Target: {}",
145 projected_column.schema(),
146 projected_schema,
147 ))), None));
148 }
149
150 let mut columns = Vec::with_capacity(projected_schema.fields.len());
151 for (idx,field) in projected_schema.fields.iter().enumerate() {
152 let column = projected_column.column(idx);
153 let extype = field.metadata().get("greptime:type").and_then(|s| ColumnExtType::from_str(s).ok());
154 let output = custom_cast(&column, field.data_type(), extype)?;
155 columns.push(output)
156 }
157 let record_batch = DfRecordBatch::try_new(projected_schema, columns)?;
158 let record_batch = if let Some(predicate) = predicate {
159 batch_filter(&record_batch, &predicate)?
160 } else {
161 record_batch
162 };
163 Ok(record_batch)
164 })
165 });
166
167 Poll::Ready(batch)
168 }
169
170 #[inline]
171 fn size_hint(&self) -> (usize, Option<usize>) {
172 self.stream.size_hint()
173 }
174}
175
176pub struct DfRecordBatchStreamAdapter {
179 stream: SendableRecordBatchStream,
180}
181
182impl DfRecordBatchStreamAdapter {
183 pub fn new(stream: SendableRecordBatchStream) -> Self {
184 Self { stream }
185 }
186}
187
188impl DfRecordBatchStream for DfRecordBatchStreamAdapter {
189 fn schema(&self) -> DfSchemaRef {
190 self.stream.schema().arrow_schema().clone()
191 }
192}
193
194impl Stream for DfRecordBatchStreamAdapter {
195 type Item = DfResult<DfRecordBatch>;
196
197 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
198 match Pin::new(&mut self.stream).poll_next(cx) {
199 Poll::Pending => Poll::Pending,
200 Poll::Ready(Some(recordbatch)) => match recordbatch {
201 Ok(recordbatch) => Poll::Ready(Some(Ok(recordbatch.into_df_record_batch()))),
202 Err(e) => Poll::Ready(Some(Err(DataFusionError::External(Box::new(e))))),
203 },
204 Poll::Ready(None) => Poll::Ready(None),
205 }
206 }
207
208 #[inline]
209 fn size_hint(&self) -> (usize, Option<usize>) {
210 self.stream.size_hint()
211 }
212}
213
214pub struct RecordBatchStreamAdapter {
218 schema: SchemaRef,
219 stream: DfSendableRecordBatchStream,
220 metrics: Option<BaselineMetrics>,
221 metrics_2: Metrics,
223 query_load_region_id: Option<u64>,
224 query_stat_counters: Option<RegionQueryStatCounters>,
225 explain_verbose: bool,
227 span: Span,
228}
229
230#[derive(Debug, Clone)]
232pub struct RegionQueryStatCounters {
233 pub query_cpu_time: Arc<AtomicU64>,
235 pub query_scanned_bytes: Arc<AtomicU64>,
237}
238
239enum Metrics {
241 Unavailable,
242 Unresolved(Arc<dyn ExecutionPlan>),
243 PartialResolved(Arc<dyn ExecutionPlan>, RecordBatchMetrics),
244 Resolved(RecordBatchMetrics),
245}
246
247impl RecordBatchStreamAdapter {
248 pub fn try_new(stream: DfSendableRecordBatchStream) -> Result<Self> {
249 let schema =
250 Arc::new(Schema::try_from(stream.schema()).context(error::SchemaConversionSnafu)?);
251 Ok(Self {
252 schema,
253 stream,
254 metrics: None,
255 metrics_2: Metrics::Unavailable,
256 query_load_region_id: None,
257 query_stat_counters: None,
258 explain_verbose: false,
259 span: Span::current(),
260 })
261 }
262
263 pub fn try_new_with_span(stream: DfSendableRecordBatchStream, span: Span) -> Result<Self> {
264 let schema =
265 Arc::new(Schema::try_from(stream.schema()).context(error::SchemaConversionSnafu)?);
266 let subspan = info_span!(parent: &span, "RecordBatchStreamAdapter");
267 Ok(Self {
268 schema,
269 stream,
270 metrics: None,
271 metrics_2: Metrics::Unavailable,
272 query_load_region_id: None,
273 query_stat_counters: None,
274 explain_verbose: false,
275 span: subspan,
276 })
277 }
278
279 pub fn set_metrics2(&mut self, plan: Arc<dyn ExecutionPlan>) {
280 self.metrics_2 = Metrics::Unresolved(plan)
281 }
282
283 fn record_query_stats_on_drop(&self) {
284 let Some(counters) = &self.query_stat_counters else {
285 return;
286 };
287
288 match &self.metrics_2 {
289 Metrics::Unresolved(df_plan) => {
290 let metrics = collect_lightweight_query_load_metrics(
291 df_plan.as_ref(),
292 self.query_load_region_id,
293 );
294 record_query_stats(counters, &metrics);
295 }
296 Metrics::PartialResolved(_, metrics) | Metrics::Resolved(metrics) => {
297 record_query_stats(counters, metrics);
298 }
299 Metrics::Unavailable => {}
300 }
301 }
302
303 pub fn set_query_load_region_id(&mut self, region_id: Option<u64>) {
304 self.query_load_region_id = region_id;
305 }
306
307 pub fn set_query_stat_counters(&mut self, counters: Option<RegionQueryStatCounters>) {
308 self.query_stat_counters = counters;
309 }
310
311 pub fn set_explain_verbose(&mut self, verbose: bool) {
313 self.explain_verbose = verbose;
314 }
315
316 fn collect_plan_metrics(&self, df_plan: &Arc<dyn ExecutionPlan>) -> RecordBatchMetrics {
317 collect_full_metrics(
318 df_plan.as_ref(),
319 self.explain_verbose,
320 self.query_load_region_id,
321 )
322 }
323
324 fn collect_partial_metrics(
325 df_plan: &dyn ExecutionPlan,
326 explain_verbose: bool,
327 query_load_region_id: Option<u64>,
328 ) -> RecordBatchMetrics {
329 if explain_verbose {
330 collect_full_metrics(df_plan, false, query_load_region_id)
333 } else {
334 collect_lightweight_query_load_metrics(df_plan, query_load_region_id)
335 }
336 }
337
338 fn update_plan_metrics(&mut self, final_metrics: bool) {
339 if final_metrics {
340 let df_plan = match &self.metrics_2 {
341 Metrics::Unresolved(df_plan) | Metrics::PartialResolved(df_plan, _) => {
342 df_plan.clone()
343 }
344 Metrics::Unavailable | Metrics::Resolved(_) => return,
345 };
346 let metrics = self.collect_plan_metrics(&df_plan);
347 self.metrics_2 = Metrics::Resolved(metrics);
348 } else {
349 let explain_verbose = self.explain_verbose;
350 let query_load_region_id = self.query_load_region_id;
351 match &mut self.metrics_2 {
352 Metrics::Unresolved(df_plan) => {
353 let df_plan = df_plan.clone();
354 let metrics = Self::collect_partial_metrics(
355 df_plan.as_ref(),
356 explain_verbose,
357 query_load_region_id,
358 );
359 self.metrics_2 = Metrics::PartialResolved(df_plan, metrics);
360 }
361 Metrics::PartialResolved(df_plan, metrics) => {
362 *metrics = Self::collect_partial_metrics(
363 df_plan.as_ref(),
364 explain_verbose,
365 query_load_region_id,
366 );
367 }
368 Metrics::Unavailable | Metrics::Resolved(_) => {}
369 }
370 }
371 }
372}
373
374pub fn region_scan_output_bytes(metrics: &RecordBatchMetrics) -> usize {
376 metrics
377 .plan_metrics
378 .iter()
379 .filter(|pm| pm.plan_name == REGION_SCAN_EXEC_NAME)
380 .flat_map(|pm| &pm.metrics)
381 .filter_map(|(name, value)| (name == "output_bytes").then_some(*value))
382 .sum()
383}
384
385fn record_query_stats(counters: &RegionQueryStatCounters, metrics: &RecordBatchMetrics) {
386 counters
387 .query_cpu_time
388 .fetch_add(metrics.elapsed_compute as u64, Ordering::Relaxed);
389 counters
390 .query_scanned_bytes
391 .fetch_add(region_scan_output_bytes(metrics) as u64, Ordering::Relaxed);
392}
393
394fn collect_full_metrics(
396 df_plan: &dyn ExecutionPlan,
397 explain_verbose: bool,
398 query_load_region_id: Option<u64>,
399) -> RecordBatchMetrics {
400 let mut metric_collector = MetricCollector::new(explain_verbose);
401 accept(df_plan, &mut metric_collector).unwrap();
402 metric_collector.record_batch_metrics.query_load_region_id = query_load_region_id;
403 metric_collector.record_batch_metrics
404}
405
406fn collect_lightweight_query_load_metrics(
413 df_plan: &dyn ExecutionPlan,
414 query_load_region_id: Option<u64>,
415) -> RecordBatchMetrics {
416 let mut metrics = RecordBatchMetrics {
417 query_load_region_id,
418 ..Default::default()
419 };
420 collect_lightweight_query_load_metrics_inner(df_plan, 0, &mut metrics);
421 metrics
422}
423
424fn collect_lightweight_query_load_metrics_inner(
427 df_plan: &dyn ExecutionPlan,
428 level: usize,
429 record_batch_metrics: &mut RecordBatchMetrics,
430) {
431 let is_region_scan = df_plan.name() == REGION_SCAN_EXEC_NAME;
432 let mut region_scan_output_bytes = None;
433
434 if let Some(metrics) = df_plan.metrics() {
435 for metric in metrics.iter() {
436 let value = metric.value();
437 match value {
438 MetricValue::ElapsedCompute(elapsed_compute) => {
439 record_batch_metrics.elapsed_compute += elapsed_compute.value();
440 }
441 MetricValue::CurrentMemoryUsage(memory_usage) => {
442 record_batch_metrics.memory_usage += memory_usage.value();
443 }
444 _ => {}
445 }
446
447 if is_region_scan && value.name() == "output_bytes" {
448 *region_scan_output_bytes.get_or_insert(0) += value.as_usize();
449 }
450 }
451 }
452
453 if let Some(output_bytes) = region_scan_output_bytes {
454 record_batch_metrics.plan_metrics.push(PlanMetrics {
455 plan: df_plan.name().to_string(),
456 plan_name: df_plan.name().to_string(),
457 level,
458 metrics: vec![("output_bytes".to_string(), output_bytes)],
459 });
460 }
461
462 for child in df_plan.children() {
463 collect_lightweight_query_load_metrics_inner(
464 child.as_ref(),
465 level + 1,
466 record_batch_metrics,
467 );
468 }
469}
470
471impl RecordBatchStream for RecordBatchStreamAdapter {
472 fn name(&self) -> &str {
473 "RecordBatchStreamAdapter"
474 }
475
476 fn schema(&self) -> SchemaRef {
477 self.schema.clone()
478 }
479
480 fn metrics(&self) -> Option<RecordBatchMetrics> {
481 match &self.metrics_2 {
482 Metrics::Unresolved(df_plan) => {
483 if self.explain_verbose {
484 Some(Self::collect_partial_metrics(
485 df_plan.as_ref(),
486 true,
487 self.query_load_region_id,
488 ))
489 } else {
490 None
491 }
492 }
493 Metrics::PartialResolved(df_plan, metrics) => Some(if self.explain_verbose {
494 Self::collect_partial_metrics(df_plan.as_ref(), true, self.query_load_region_id)
495 } else {
496 metrics.clone()
497 }),
498 Metrics::Resolved(metrics) => Some(metrics.clone()),
499 Metrics::Unavailable => None,
500 }
501 }
502
503 fn output_ordering(&self) -> Option<&[OrderOption]> {
504 None
505 }
506}
507
508impl Stream for RecordBatchStreamAdapter {
509 type Item = Result<RecordBatch>;
510
511 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
512 let timer = self
513 .metrics
514 .as_ref()
515 .map(|m| m.elapsed_compute().clone())
516 .unwrap_or_default();
517 let _guard = timer.timer();
518 let poll_span = info_span!(parent: &self.span, "poll_next");
519 let _entered = poll_span.enter();
520 match Pin::new(&mut self.stream).poll_next(cx) {
521 Poll::Pending => Poll::Pending,
522 Poll::Ready(Some(df_record_batch)) => {
523 let df_record_batch = df_record_batch?;
524 self.update_plan_metrics(false);
525 Poll::Ready(Some(Ok(RecordBatch::from_df_record_batch(
526 self.schema(),
527 df_record_batch,
528 ))))
529 }
530 Poll::Ready(None) => {
531 self.update_plan_metrics(true);
532 Poll::Ready(None)
533 }
534 }
535 }
536
537 #[inline]
538 fn size_hint(&self) -> (usize, Option<usize>) {
539 self.stream.size_hint()
540 }
541}
542
543impl Drop for RecordBatchStreamAdapter {
544 fn drop(&mut self) {
545 self.record_query_stats_on_drop();
546 }
547}
548
549pub struct MetricCollector {
551 current_level: usize,
552 pub record_batch_metrics: RecordBatchMetrics,
553 verbose: bool,
554}
555
556impl MetricCollector {
557 pub fn new(verbose: bool) -> Self {
558 Self {
559 current_level: 0,
560 record_batch_metrics: RecordBatchMetrics::default(),
561 verbose,
562 }
563 }
564}
565
566impl ExecutionPlanVisitor for MetricCollector {
567 type Error = !;
568
569 fn pre_visit(&mut self, plan: &dyn ExecutionPlan) -> std::result::Result<bool, Self::Error> {
570 let Some(metric) = plan.metrics() else {
572 self.record_batch_metrics.plan_metrics.push(PlanMetrics {
573 plan: plan.name().to_string(),
574 plan_name: plan.name().to_string(),
575 level: self.current_level,
576 metrics: vec![],
577 });
578 self.current_level += 1;
579 return Ok(true);
580 };
581
582 let metric = metric
584 .aggregate_by_name()
585 .sorted_for_display()
586 .timestamps_removed();
587 let mut plan_metric = PlanMetrics {
588 plan: one_line(plan, self.verbose).to_string(),
589 plan_name: plan.name().to_string(),
590 level: self.current_level,
591 metrics: Vec::with_capacity(metric.iter().size_hint().0),
592 };
593 for m in metric.iter() {
594 plan_metric
595 .metrics
596 .push((m.value().name().to_string(), m.value().as_usize()));
597
598 match m.value() {
600 MetricValue::ElapsedCompute(ec) => {
601 self.record_batch_metrics.elapsed_compute += ec.value()
602 }
603 MetricValue::CurrentMemoryUsage(m) => {
604 self.record_batch_metrics.memory_usage += m.value()
605 }
606 _ => {}
607 }
608 }
609 self.record_batch_metrics.plan_metrics.push(plan_metric);
610
611 self.current_level += 1;
612 Ok(true)
613 }
614
615 fn post_visit(&mut self, _plan: &dyn ExecutionPlan) -> std::result::Result<bool, Self::Error> {
616 self.current_level -= 1;
617 Ok(true)
618 }
619}
620
621fn one_line(plan: &dyn ExecutionPlan, verbose: bool) -> impl fmt::Display + '_ {
624 struct Wrapper<'a> {
625 plan: &'a dyn ExecutionPlan,
626 format_type: DisplayFormatType,
627 }
628
629 impl fmt::Display for Wrapper<'_> {
630 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
631 self.plan.fmt_as(self.format_type, f)?;
632 writeln!(f)
633 }
634 }
635
636 let format_type = if verbose {
637 DisplayFormatType::Verbose
638 } else {
639 DisplayFormatType::Default
640 };
641 Wrapper { plan, format_type }
642}
643
644#[derive(serde::Serialize, serde::Deserialize, Default, Debug, Clone)]
647pub struct RecordBatchMetrics {
648 pub elapsed_compute: usize,
651 pub memory_usage: usize,
653 pub plan_metrics: Vec<PlanMetrics>,
656 #[serde(default, skip_serializing_if = "Option::is_none")]
658 pub query_load_region_id: Option<u64>,
659 #[serde(default, skip_serializing_if = "Vec::is_empty")]
672 pub region_watermarks: Vec<RegionWatermarkEntry>,
673}
674
675#[derive(serde::Serialize, serde::Deserialize, Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
676pub struct RegionWatermarkEntry {
677 pub region_id: u64,
678 #[serde(default, skip_serializing_if = "Option::is_none")]
679 pub watermark: Option<u64>,
680}
681
682fn is_time_metric(metric_name: &str) -> bool {
684 metric_name.contains("elapsed") || metric_name.contains("time") || metric_name.contains("cost")
685}
686
687fn is_bytes_metric(metric_name: &str) -> bool {
689 metric_name.contains("bytes") || metric_name.contains("mem")
690}
691
692fn format_bytes_human_readable(bytes: usize) -> String {
693 format!("{}", ReadableSize(bytes as u64))
694}
695
696impl Display for RecordBatchMetrics {
698 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
699 for metric in &self.plan_metrics {
700 write!(
701 f,
702 "{:indent$}{} metrics=[",
703 " ",
704 metric.plan.trim_end(),
705 indent = metric.level * 2,
706 )?;
707 for (label, value) in &metric.metrics {
708 if is_time_metric(label) {
709 write!(
710 f,
711 "{}: {}, ",
712 label,
713 format_nanoseconds_human_readable(*value),
714 )?;
715 } else if is_bytes_metric(label) {
716 write!(f, "{}: {}, ", label, format_bytes_human_readable(*value),)?;
717 } else {
718 write!(f, "{}: {}, ", label, value)?;
719 }
720 }
721 writeln!(f, "]")?;
722 }
723
724 Ok(())
725 }
726}
727
728#[derive(serde::Serialize, serde::Deserialize, Default, Debug, Clone)]
729pub struct PlanMetrics {
730 pub plan: String,
732 #[serde(default)]
734 pub plan_name: String,
735 pub level: usize,
737 pub metrics: Vec<(String, usize)>,
740}
741
742enum AsyncRecordBatchStreamAdapterState {
743 Uninit(FutureStream),
744 Ready(SendableRecordBatchStream),
745 Failed,
746}
747
748pub struct AsyncRecordBatchStreamAdapter {
749 schema: SchemaRef,
750 state: AsyncRecordBatchStreamAdapterState,
751}
752
753impl AsyncRecordBatchStreamAdapter {
754 pub fn new(schema: SchemaRef, stream: FutureStream) -> Self {
755 Self {
756 schema,
757 state: AsyncRecordBatchStreamAdapterState::Uninit(stream),
758 }
759 }
760}
761
762impl RecordBatchStream for AsyncRecordBatchStreamAdapter {
763 fn schema(&self) -> SchemaRef {
764 self.schema.clone()
765 }
766
767 fn output_ordering(&self) -> Option<&[OrderOption]> {
768 None
769 }
770
771 fn metrics(&self) -> Option<RecordBatchMetrics> {
772 None
773 }
774}
775
776impl Stream for AsyncRecordBatchStreamAdapter {
777 type Item = Result<RecordBatch>;
778
779 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
780 loop {
781 match &mut self.state {
782 AsyncRecordBatchStreamAdapterState::Uninit(stream_future) => {
783 match ready!(Pin::new(stream_future).poll(cx)) {
784 Ok(stream) => {
785 self.state = AsyncRecordBatchStreamAdapterState::Ready(stream);
786 continue;
787 }
788 Err(e) => {
789 self.state = AsyncRecordBatchStreamAdapterState::Failed;
790 return Poll::Ready(Some(Err(e)));
791 }
792 };
793 }
794 AsyncRecordBatchStreamAdapterState::Ready(stream) => {
795 return Poll::Ready(ready!(Pin::new(stream).poll_next(cx)));
796 }
797 AsyncRecordBatchStreamAdapterState::Failed => return Poll::Ready(None),
798 }
799 }
800 }
801
802 #[inline]
804 fn size_hint(&self) -> (usize, Option<usize>) {
805 (0, None)
806 }
807}
808
809fn custom_cast(
811 array: &dyn Array,
812 target_type: &ArrowDataType,
813 extype: Option<ColumnExtType>,
814) -> std::result::Result<Arc<dyn Array>, ArrowError> {
815 if let ArrowDataType::Map(_, _) = array.data_type()
816 && let ArrowDataType::Binary = target_type
817 {
818 return convert_map_to_json_binary(array, extype);
819 }
820
821 cast(array, target_type)
822}
823
824fn convert_map_to_json_binary(
826 array: &dyn Array,
827 extype: Option<ColumnExtType>,
828) -> std::result::Result<Arc<dyn Array>, ArrowError> {
829 use datatypes::arrow::array::{BinaryArray, MapArray};
830 use serde_json::Value;
831
832 let map_array = array
833 .as_any()
834 .downcast_ref::<MapArray>()
835 .ok_or_else(|| ArrowError::CastError("Failed to downcast to MapArray".to_string()))?;
836
837 let mut json_values = Vec::with_capacity(map_array.len());
838
839 for i in 0..map_array.len() {
840 if map_array.is_null(i) {
841 json_values.push(None);
842 } else {
843 let map_entry = map_array.value(i);
845 let key_value_array = map_entry
846 .as_any()
847 .downcast_ref::<datatypes::arrow::array::StructArray>()
848 .ok_or_else(|| {
849 ArrowError::CastError("Failed to downcast to StructArray".to_string())
850 })?;
851
852 let mut json_obj = serde_json::Map::with_capacity(key_value_array.len());
854
855 for j in 0..key_value_array.len() {
856 if key_value_array.is_null(j) {
857 continue;
858 }
859 let key_field = key_value_array.column(0);
860 let value_field = key_value_array.column(1);
861
862 if key_field.is_null(j) {
863 continue;
864 }
865
866 let key = key_field
867 .as_any()
868 .downcast_ref::<datatypes::arrow::array::StringArray>()
869 .ok_or_else(|| {
870 ArrowError::CastError("Failed to downcast key to StringArray".to_string())
871 })?
872 .value(j);
873
874 let value = if value_field.is_null(j) {
875 Value::Null
876 } else {
877 let value_str = value_field
878 .as_any()
879 .downcast_ref::<datatypes::arrow::array::StringArray>()
880 .ok_or_else(|| {
881 ArrowError::CastError(
882 "Failed to downcast value to StringArray".to_string(),
883 )
884 })?
885 .value(j);
886 Value::String(value_str.to_string())
887 };
888
889 json_obj.insert(key.to_string(), value);
890 }
891
892 let json_value = Value::Object(json_obj);
893 let json_bytes = match extype {
894 Some(ColumnExtType::Json) => {
895 let json_string = match serde_json::to_string(&json_value) {
896 Ok(s) => s,
897 Err(e) => {
898 return Err(ArrowError::CastError(format!(
899 "Failed to serialize JSON: {}",
900 e
901 )));
902 }
903 };
904 match jsonb::parse_value(json_string.as_bytes()) {
905 Ok(jsonb_value) => jsonb_value.to_vec(),
906 Err(e) => {
907 return Err(ArrowError::CastError(format!(
908 "Failed to serialize JSONB: {}",
909 e
910 )));
911 }
912 }
913 }
914 _ => match serde_json::to_vec(&json_value) {
915 Ok(b) => b,
916 Err(e) => {
917 return Err(ArrowError::CastError(format!(
918 "Failed to serialize JSON: {}",
919 e
920 )));
921 }
922 },
923 };
924 json_values.push(Some(json_bytes));
925 }
926 }
927
928 let binary_array = BinaryArray::from_iter(json_values);
929 Ok(Arc::new(binary_array))
930}
931
932#[cfg(test)]
933mod test {
934 use std::any::Any;
935 use std::time::Duration;
936
937 use common_error::ext::BoxedError;
938 use common_error::mock::MockError;
939 use common_error::status_code::StatusCode;
940 use datafusion::execution::TaskContext;
941 use datafusion::physical_expr::{EquivalenceProperties, Partitioning};
942 use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
943 use datafusion::physical_plan::metrics::{ExecutionPlanMetricsSet, MetricBuilder, MetricsSet};
944 use datafusion::physical_plan::{DisplayAs, PlanProperties};
945 use datatypes::arrow::array::{ArrayRef, MapArray, StringArray, StructArray};
946 use datatypes::arrow::buffer::OffsetBuffer;
947 use datatypes::arrow::datatypes::Field;
948 use datatypes::prelude::ConcreteDataType;
949 use datatypes::schema::ColumnSchema;
950 use datatypes::vectors::Int32Vector;
951 use futures::StreamExt;
952 use serde_json::json;
953 use snafu::IntoError;
954
955 use super::*;
956 use crate::RecordBatches;
957 use crate::error::Error;
958
959 #[derive(Debug)]
960 struct TestMetricsExec {
961 properties: Arc<PlanProperties>,
962 metrics: ExecutionPlanMetricsSet,
963 format_plan: bool,
964 }
965
966 impl TestMetricsExec {
967 fn new(schema: DfSchemaRef) -> Self {
968 Self::with_output_bytes(schema, &[24])
969 }
970
971 fn new_without_plan_formatting(schema: DfSchemaRef) -> Self {
972 Self::with_output_bytes_and_formatting(schema, &[24], false)
973 }
974
975 fn with_output_bytes(schema: DfSchemaRef, output_bytes_by_partition: &[usize]) -> Self {
976 Self::with_output_bytes_and_formatting(schema, output_bytes_by_partition, true)
977 }
978
979 fn with_output_bytes_and_formatting(
980 schema: DfSchemaRef,
981 output_bytes_by_partition: &[usize],
982 format_plan: bool,
983 ) -> Self {
984 let metrics = ExecutionPlanMetricsSet::new();
985 let elapsed_compute = MetricBuilder::new(&metrics).elapsed_compute(0);
986 elapsed_compute.add_duration(Duration::from_nanos(42));
987 for (partition, output_bytes) in output_bytes_by_partition.iter().copied().enumerate() {
988 let metric = MetricBuilder::new(&metrics).output_bytes(partition);
989 metric.add(output_bytes);
990 }
991
992 Self {
993 properties: Arc::new(PlanProperties::new(
994 EquivalenceProperties::new(schema),
995 Partitioning::UnknownPartitioning(output_bytes_by_partition.len().max(1)),
996 EmissionType::Incremental,
997 Boundedness::Bounded,
998 )),
999 metrics,
1000 format_plan,
1001 }
1002 }
1003 }
1004
1005 impl DisplayAs for TestMetricsExec {
1006 fn fmt_as(&self, t: DisplayFormatType, f: &mut std::fmt::Formatter) -> std::fmt::Result {
1007 assert!(
1008 self.format_plan,
1009 "non-verbose lightweight partial metrics must not format the plan"
1010 );
1011 write!(f, "RegionScanExec")?;
1012 if matches!(t, DisplayFormatType::Verbose) {
1013 write!(f, ": files=[file-1.parquet]")?;
1014 }
1015 Ok(())
1016 }
1017 }
1018
1019 impl ExecutionPlan for TestMetricsExec {
1020 fn name(&self) -> &str {
1021 REGION_SCAN_EXEC_NAME
1022 }
1023
1024 fn as_any(&self) -> &dyn Any {
1025 self
1026 }
1027
1028 fn properties(&self) -> &Arc<PlanProperties> {
1029 &self.properties
1030 }
1031
1032 fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
1033 vec![]
1034 }
1035
1036 fn with_new_children(
1037 self: Arc<Self>,
1038 _children: Vec<Arc<dyn ExecutionPlan>>,
1039 ) -> datafusion_common::Result<Arc<dyn ExecutionPlan>> {
1040 Ok(self)
1041 }
1042
1043 fn execute(
1044 &self,
1045 _partition: usize,
1046 _context: Arc<TaskContext>,
1047 ) -> datafusion_common::Result<DfSendableRecordBatchStream> {
1048 unreachable!("the test passes a separate stream to RecordBatchStreamAdapter")
1049 }
1050
1051 fn metrics(&self) -> Option<MetricsSet> {
1052 Some(self.metrics.clone_inner())
1053 }
1054 }
1055
1056 #[test]
1057 fn test_lightweight_query_load_metrics_sums_region_scan_output_bytes() {
1058 let schema = Arc::new(Schema::new(vec![ColumnSchema::new(
1059 "a",
1060 ConcreteDataType::int32_datatype(),
1061 false,
1062 )]));
1063 let plan = TestMetricsExec::with_output_bytes(schema.arrow_schema().clone(), &[24, 18]);
1064
1065 let metrics = collect_lightweight_query_load_metrics(&plan, Some(42));
1066
1067 assert_eq!(metrics.query_load_region_id, Some(42));
1068 assert_eq!(region_scan_output_bytes(&metrics), 42);
1069 assert_eq!(metrics.plan_metrics.len(), 1);
1070 assert_eq!(
1071 metrics.plan_metrics[0].metrics,
1072 vec![("output_bytes".to_string(), 42)]
1073 );
1074 }
1075
1076 #[tokio::test]
1077 async fn test_record_batch_stream_adapter_collects_lightweight_partial_metrics() {
1078 let schema = Arc::new(Schema::new(vec![ColumnSchema::new(
1079 "a",
1080 ConcreteDataType::int32_datatype(),
1081 false,
1082 )]));
1083 let batch1 = RecordBatch::new(
1084 schema.clone(),
1085 vec![Arc::new(Int32Vector::from_slice([1])) as _],
1086 )
1087 .unwrap()
1088 .into_df_record_batch();
1089 let batch2 = RecordBatch::new(
1090 schema.clone(),
1091 vec![Arc::new(Int32Vector::from_slice([2])) as _],
1092 )
1093 .unwrap()
1094 .into_df_record_batch();
1095 let df_stream = Box::pin(
1096 datafusion::physical_plan::stream::RecordBatchStreamAdapter::new(
1097 schema.arrow_schema().clone(),
1098 futures::stream::iter(vec![Ok(batch1), Ok(batch2)]),
1099 ),
1100 );
1101 let plan = Arc::new(TestMetricsExec::new_without_plan_formatting(
1102 schema.arrow_schema().clone(),
1103 ));
1104
1105 let mut adapter = RecordBatchStreamAdapter::try_new(df_stream).unwrap();
1106 adapter.set_metrics2(plan);
1107 adapter.set_query_load_region_id(Some(42));
1108
1109 assert!(adapter.metrics().is_none());
1110 assert!(adapter.next().await.unwrap().is_ok());
1111 let metrics = adapter
1112 .metrics()
1113 .expect("non-verbose queries need partial query-load metrics before EOF");
1114 assert_eq!(metrics.elapsed_compute, 42);
1115 assert_eq!(metrics.query_load_region_id, Some(42));
1116 assert_eq!(region_scan_output_bytes(&metrics), 24);
1117 assert_eq!(metrics.plan_metrics.len(), 1);
1118 assert_eq!(metrics.plan_metrics[0].plan, REGION_SCAN_EXEC_NAME);
1119 assert_eq!(metrics.plan_metrics[0].plan_name, REGION_SCAN_EXEC_NAME);
1120 assert!(
1121 metrics.plan_metrics[0]
1122 .metrics
1123 .iter()
1124 .any(|(name, value)| name == "output_bytes" && *value == 24)
1125 );
1126 }
1127
1128 #[tokio::test]
1129 async fn test_record_batch_stream_adapter_uses_compact_partial_and_verbose_final_metrics() {
1130 let schema = Arc::new(Schema::new(vec![ColumnSchema::new(
1131 "a",
1132 ConcreteDataType::int32_datatype(),
1133 false,
1134 )]));
1135 let batch = RecordBatch::new(
1136 schema.clone(),
1137 vec![Arc::new(Int32Vector::from_slice([1])) as _],
1138 )
1139 .unwrap()
1140 .into_df_record_batch();
1141 let df_stream = Box::pin(
1142 datafusion::physical_plan::stream::RecordBatchStreamAdapter::new(
1143 schema.arrow_schema().clone(),
1144 futures::stream::iter(vec![Ok(batch)]),
1145 ),
1146 );
1147 let plan = Arc::new(TestMetricsExec::new(schema.arrow_schema().clone()));
1148 let mut adapter = RecordBatchStreamAdapter::try_new(df_stream).unwrap();
1149 adapter.set_metrics2(plan);
1150 adapter.set_explain_verbose(true);
1151
1152 adapter.next().await.unwrap().unwrap();
1153 let partial = adapter.metrics().unwrap();
1154 assert_eq!(partial.plan_metrics.len(), 1);
1155 assert_eq!(partial.plan_metrics[0].plan.trim_end(), "RegionScanExec");
1156 assert!(!partial.plan_metrics[0].plan.contains("files"));
1157
1158 assert!(adapter.next().await.is_none());
1159 let final_metrics = adapter.metrics().unwrap();
1160 assert!(final_metrics.plan_metrics[0].plan.contains("files"));
1161 }
1162
1163 #[test]
1164 fn test_record_batch_stream_adapter_reuses_partial_query_stats_on_drop() {
1165 let schema = Arc::new(Schema::new(vec![ColumnSchema::new(
1166 "a",
1167 ConcreteDataType::int32_datatype(),
1168 false,
1169 )]));
1170 let df_stream = Box::pin(
1171 datafusion::physical_plan::stream::RecordBatchStreamAdapter::new(
1172 schema.arrow_schema().clone(),
1173 futures::stream::empty::<datafusion::error::Result<DfRecordBatch>>(),
1174 ),
1175 );
1176 let counters = RegionQueryStatCounters {
1177 query_cpu_time: Arc::new(AtomicU64::new(10)),
1178 query_scanned_bytes: Arc::new(AtomicU64::new(20)),
1179 };
1180 let stale_metrics = RecordBatchMetrics {
1181 elapsed_compute: 1,
1182 plan_metrics: vec![PlanMetrics {
1183 plan: REGION_SCAN_EXEC_NAME.to_string(),
1184 plan_name: REGION_SCAN_EXEC_NAME.to_string(),
1185 level: 0,
1186 metrics: vec![("output_bytes".to_string(), 2)],
1187 }],
1188 ..Default::default()
1189 };
1190 let adapter = RecordBatchStreamAdapter {
1191 schema: schema.clone(),
1192 stream: df_stream,
1193 metrics: None,
1194 metrics_2: Metrics::PartialResolved(
1195 Arc::new(TestMetricsExec::new(schema.arrow_schema().clone())),
1196 stale_metrics,
1197 ),
1198 query_load_region_id: None,
1199 query_stat_counters: Some(counters.clone()),
1200 explain_verbose: false,
1201 span: Span::current(),
1202 };
1203
1204 drop(adapter);
1205
1206 assert_eq!(counters.query_cpu_time.load(Ordering::Relaxed), 11);
1207 assert_eq!(counters.query_scanned_bytes.load(Ordering::Relaxed), 22);
1208 }
1209
1210 #[tokio::test]
1211 async fn test_async_recordbatch_stream_adaptor() {
1212 struct MaybeErrorRecordBatchStream {
1213 items: Vec<Result<RecordBatch>>,
1214 }
1215
1216 impl RecordBatchStream for MaybeErrorRecordBatchStream {
1217 fn schema(&self) -> SchemaRef {
1218 unimplemented!()
1219 }
1220
1221 fn output_ordering(&self) -> Option<&[OrderOption]> {
1222 None
1223 }
1224
1225 fn metrics(&self) -> Option<RecordBatchMetrics> {
1226 None
1227 }
1228 }
1229
1230 impl Stream for MaybeErrorRecordBatchStream {
1231 type Item = Result<RecordBatch>;
1232
1233 fn poll_next(
1234 mut self: Pin<&mut Self>,
1235 _: &mut Context<'_>,
1236 ) -> Poll<Option<Self::Item>> {
1237 if let Some(batch) = self.items.pop() {
1238 Poll::Ready(Some(Ok(batch?)))
1239 } else {
1240 Poll::Ready(None)
1241 }
1242 }
1243 }
1244
1245 fn new_future_stream(
1246 maybe_recordbatches: Result<Vec<Result<RecordBatch>>>,
1247 ) -> FutureStream {
1248 Box::pin(async move {
1249 maybe_recordbatches
1250 .map(|items| Box::pin(MaybeErrorRecordBatchStream { items }) as _)
1251 })
1252 }
1253
1254 let schema = Arc::new(Schema::new(vec![ColumnSchema::new(
1255 "a",
1256 ConcreteDataType::int32_datatype(),
1257 false,
1258 )]));
1259 let batch1 = RecordBatch::new(
1260 schema.clone(),
1261 vec![Arc::new(Int32Vector::from_slice([1])) as _],
1262 )
1263 .unwrap();
1264 let batch2 = RecordBatch::new(
1265 schema.clone(),
1266 vec![Arc::new(Int32Vector::from_slice([2])) as _],
1267 )
1268 .unwrap();
1269
1270 let success_stream = new_future_stream(Ok(vec![Ok(batch1.clone()), Ok(batch2.clone())]));
1271 let adapter = AsyncRecordBatchStreamAdapter::new(schema.clone(), success_stream);
1272 let collected = RecordBatches::try_collect(Box::pin(adapter)).await.unwrap();
1273 assert_eq!(
1274 collected,
1275 RecordBatches::try_new(schema.clone(), vec![batch2.clone(), batch1.clone()]).unwrap()
1276 );
1277
1278 let poll_err_stream = new_future_stream(Ok(vec![
1279 Ok(batch1.clone()),
1280 Err(error::ExternalSnafu
1281 .into_error(BoxedError::new(MockError::new(StatusCode::Unknown)))),
1282 ]));
1283 let adapter = AsyncRecordBatchStreamAdapter::new(schema.clone(), poll_err_stream);
1284 let err = RecordBatches::try_collect(Box::pin(adapter))
1285 .await
1286 .unwrap_err();
1287 assert!(
1288 matches!(err, Error::External { .. }),
1289 "unexpected err {err}"
1290 );
1291
1292 let failed_to_init_stream =
1293 new_future_stream(Err(error::ExternalSnafu
1294 .into_error(BoxedError::new(MockError::new(StatusCode::Internal)))));
1295 let adapter = AsyncRecordBatchStreamAdapter::new(schema.clone(), failed_to_init_stream);
1296 let err = RecordBatches::try_collect(Box::pin(adapter))
1297 .await
1298 .unwrap_err();
1299 assert!(
1300 matches!(err, Error::External { .. }),
1301 "unexpected err {err}"
1302 );
1303 }
1304
1305 #[test]
1306 fn test_convert_map_to_json_binary() {
1307 let keys = StringArray::from(vec![Some("a"), Some("b"), Some("c"), Some("x")]);
1308 let values = StringArray::from(vec![Some("1"), None, Some("3"), Some("42")]);
1309 let key_field = Arc::new(Field::new("key", ArrowDataType::Utf8, false));
1310 let value_field = Arc::new(Field::new("value", ArrowDataType::Utf8, true));
1311 let struct_type = ArrowDataType::Struct(vec![key_field, value_field].into());
1312
1313 let entries_field = Arc::new(Field::new("entries", struct_type, false));
1314
1315 let struct_array = StructArray::from(vec![
1316 (
1317 Arc::new(Field::new("key", ArrowDataType::Utf8, false)),
1318 Arc::new(keys) as ArrayRef,
1319 ),
1320 (
1321 Arc::new(Field::new("value", ArrowDataType::Utf8, true)),
1322 Arc::new(values) as ArrayRef,
1323 ),
1324 ]);
1325
1326 let offsets = OffsetBuffer::from_lengths([3, 0, 1]);
1327 let nulls = datatypes::arrow::buffer::NullBuffer::from(vec![true, false, true]);
1328
1329 let map_array = MapArray::new(
1330 entries_field,
1331 offsets,
1332 struct_array,
1333 Some(nulls), false,
1335 );
1336
1337 let result = convert_map_to_json_binary(&map_array, None).unwrap();
1338 let binary_array = result
1339 .as_any()
1340 .downcast_ref::<datatypes::arrow::array::BinaryArray>()
1341 .unwrap();
1342
1343 let expected_jsons = [
1344 Some(r#"{"a":"1","b":null,"c":"3"}"#),
1345 None,
1346 Some(r#"{"x":"42"}"#),
1347 ];
1348
1349 for (i, _) in expected_jsons.iter().enumerate() {
1350 if let Some(expected) = &expected_jsons[i] {
1351 assert!(!binary_array.is_null(i));
1352 let actual_bytes = binary_array.value(i);
1353 let actual_str = std::str::from_utf8(actual_bytes).unwrap();
1354 assert_eq!(actual_str, *expected);
1355 } else {
1356 assert!(binary_array.is_null(i));
1357 }
1358 }
1359
1360 let result_json =
1361 convert_map_to_json_binary(&map_array, Some(ColumnExtType::Json)).unwrap();
1362 let binary_array_json = result_json
1363 .as_any()
1364 .downcast_ref::<datatypes::arrow::array::BinaryArray>()
1365 .unwrap();
1366
1367 for (i, _) in expected_jsons.iter().enumerate() {
1368 if expected_jsons[i].is_some() {
1369 assert!(!binary_array_json.is_null(i));
1370 let actual_bytes = binary_array_json.value(i);
1371 assert_ne!(actual_bytes, expected_jsons[i].unwrap().as_bytes());
1372 } else {
1373 assert!(binary_array_json.is_null(i));
1374 }
1375 }
1376 }
1377
1378 #[test]
1379 fn test_record_query_stats_updates_region_counters() {
1380 let counters = RegionQueryStatCounters {
1381 query_cpu_time: Arc::new(AtomicU64::new(10)),
1382 query_scanned_bytes: Arc::new(AtomicU64::new(20)),
1383 };
1384 let metrics = RecordBatchMetrics {
1385 elapsed_compute: 2_000_000,
1386 plan_metrics: vec![PlanMetrics {
1387 plan: "RegionScanExec: region=1".to_string(),
1388 plan_name: REGION_SCAN_EXEC_NAME.to_string(),
1389 level: 0,
1390 metrics: vec![("output_bytes".to_string(), 42)],
1391 }],
1392 ..Default::default()
1393 };
1394
1395 record_query_stats(&counters, &metrics);
1396
1397 assert_eq!(counters.query_cpu_time.load(Ordering::Relaxed), 2_000_010);
1398 assert_eq!(counters.query_scanned_bytes.load(Ordering::Relaxed), 62);
1399 }
1400
1401 #[test]
1402 fn test_record_batch_stream_adapter_records_query_stats_on_drop() {
1403 let schema = Arc::new(Schema::new(vec![ColumnSchema::new(
1404 "a",
1405 ConcreteDataType::int32_datatype(),
1406 false,
1407 )]));
1408 let df_stream = Box::pin(
1409 datafusion::physical_plan::stream::RecordBatchStreamAdapter::new(
1410 schema.arrow_schema().clone(),
1411 futures::stream::empty::<datafusion::error::Result<DfRecordBatch>>(),
1412 ),
1413 );
1414 let counters = RegionQueryStatCounters {
1415 query_cpu_time: Arc::new(AtomicU64::new(10)),
1416 query_scanned_bytes: Arc::new(AtomicU64::new(20)),
1417 };
1418 let metrics = RecordBatchMetrics {
1419 elapsed_compute: 2_000_000,
1420 plan_metrics: vec![PlanMetrics {
1421 plan: "RegionScanExec: region=1".to_string(),
1422 plan_name: REGION_SCAN_EXEC_NAME.to_string(),
1423 level: 0,
1424 metrics: vec![("output_bytes".to_string(), 42)],
1425 }],
1426 ..Default::default()
1427 };
1428 let adapter = RecordBatchStreamAdapter {
1429 schema,
1430 stream: df_stream,
1431 metrics: None,
1432 metrics_2: Metrics::Resolved(metrics),
1433 query_load_region_id: None,
1434 query_stat_counters: Some(counters.clone()),
1435 explain_verbose: false,
1436 span: Span::current(),
1437 };
1438
1439 drop(adapter);
1440
1441 assert_eq!(counters.query_cpu_time.load(Ordering::Relaxed), 2_000_010);
1442 assert_eq!(counters.query_scanned_bytes.load(Ordering::Relaxed), 62);
1443 }
1444
1445 #[test]
1446 fn test_recordbatch_metrics_deserializes_without_region_watermarks() {
1447 let metrics: RecordBatchMetrics = serde_json::from_value(json!({
1448 "elapsed_compute": 12,
1449 "memory_usage": 34,
1450 "plan_metrics": []
1451 }))
1452 .unwrap();
1453
1454 assert!(metrics.region_watermarks.is_empty());
1455 assert_eq!(metrics.elapsed_compute, 12);
1456 assert_eq!(metrics.memory_usage, 34);
1457 }
1458
1459 #[test]
1460 fn test_plan_metrics_deserializes_without_plan_name() {
1461 let metrics: RecordBatchMetrics = serde_json::from_value(json!({
1462 "elapsed_compute": 12,
1463 "memory_usage": 34,
1464 "plan_metrics": [{
1465 "plan": "SeqScan: region=1",
1466 "level": 0,
1467 "metrics": []
1468 }]
1469 }))
1470 .unwrap();
1471
1472 assert_eq!(metrics.plan_metrics[0].plan_name, "");
1473 }
1474
1475 #[test]
1476 fn test_recordbatch_metrics_region_watermarks_serde_roundtrip() {
1477 let metrics = RecordBatchMetrics {
1478 region_watermarks: vec![
1479 RegionWatermarkEntry {
1480 region_id: 1,
1481 watermark: Some(100),
1482 },
1483 RegionWatermarkEntry {
1484 region_id: 2,
1485 watermark: None,
1486 },
1487 ],
1488 ..Default::default()
1489 };
1490
1491 let value = serde_json::to_value(&metrics).unwrap();
1492 assert_eq!(
1493 value.get("region_watermarks").unwrap(),
1494 &json!([
1495 { "region_id": 1, "watermark": 100 },
1496 { "region_id": 2 }
1497 ])
1498 );
1499
1500 let decoded: RecordBatchMetrics = serde_json::from_value(value).unwrap();
1501 assert_eq!(decoded.region_watermarks, metrics.region_watermarks);
1502 }
1503
1504 #[test]
1505 fn test_recordbatch_metrics_skips_empty_region_watermarks_on_serialize() {
1506 let value = serde_json::to_value(RecordBatchMetrics::default()).unwrap();
1507 assert!(value.get("region_watermarks").is_none());
1508 }
1509}