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