Skip to main content

common_recordbatch/
adapter.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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/// Casts the `RecordBatch`es of `stream` against the `output_schema`.
62#[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
176/// Greptime SendableRecordBatchStream -> DataFusion RecordBatchStream.
177/// The reverse one is [RecordBatchStreamAdapter].
178pub 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
214/// DataFusion [SendableRecordBatchStream](DfSendableRecordBatchStream) -> Greptime [RecordBatchStream].
215/// The reverse one is [DfRecordBatchStreamAdapter].
216/// It can collect metrics from DataFusion execution plan.
217pub struct RecordBatchStreamAdapter {
218    schema: SchemaRef,
219    stream: DfSendableRecordBatchStream,
220    metrics: Option<BaselineMetrics>,
221    /// Aggregated plan-level metrics. Resolved after an [ExecutionPlan] is finished.
222    metrics_2: Metrics,
223    query_load_region_id: Option<u64>,
224    query_stat_counters: Option<RegionQueryStatCounters>,
225    /// Display plan and metrics in verbose mode.
226    explain_verbose: bool,
227    span: Span,
228}
229
230/// Query statistic counters owned by a region.
231#[derive(Debug, Clone)]
232pub struct RegionQueryStatCounters {
233    /// The total query CPU time in nanoseconds.
234    pub query_cpu_time: Arc<AtomicU64>,
235    /// The total scanned bytes.
236    pub query_scanned_bytes: Arc<AtomicU64>,
237}
238
239/// Json encoded metrics. Contains metric from a whole plan tree.
240enum 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    /// Set the verbose mode for displaying plan and metrics.
312    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            // Verbose in-progress snapshots have the same complete topology as
331            // the final snapshot, but use compact (Default) plan formatting.
332            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
374/// Extracts total `output_bytes` from region scan plan nodes.
375pub 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
394/// Collects the complete plan metrics used by terminal metrics and verbose analyze output.
395fn 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
406/// Collects the minimal metrics needed for query-load reporting before EOF.
407///
408/// This intentionally avoids [`MetricCollector`]'s full per-plan aggregation,
409/// sorting, and plan formatting work on the normal query hot path. The result
410/// is still enough for early-stop/cancellation paths to report per-region CPU
411/// time, scanned bytes, and physical-region attribution.
412fn 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
424/// Recursively walks the physical plan and reads raw metric values without
425/// formatting plan nodes.
426fn 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
549/// An [ExecutionPlanVisitor] to collect metrics from a [ExecutionPlan].
550pub 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        // skip if no metric available
571        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        // scrape plan metrics
583        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            // aggregate high-level metrics
599            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
621/// Returns a single-line summary of the root of the plan.
622/// If the `verbose` flag is set, it will display detailed information about the plan.
623fn 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/// [`RecordBatchMetrics`] carrys metrics value
645/// from datanode to frontend through gRPC
646#[derive(serde::Serialize, serde::Deserialize, Default, Debug, Clone)]
647pub struct RecordBatchMetrics {
648    // High-level aggregated metrics
649    /// CPU consumption in nanoseconds
650    pub elapsed_compute: usize,
651    /// Memory used by the plan in bytes
652    pub memory_usage: usize,
653    // Detailed per-plan metrics
654    /// An ordered list of plan metrics, from top to bottom in post-order.
655    pub plan_metrics: Vec<PlanMetrics>,
656    /// Region id that should receive query-load metrics for this scan.
657    #[serde(default, skip_serializing_if = "Option::is_none")]
658    pub query_load_region_id: Option<u64>,
659    /// Per-region watermark for incremental-read checkpoint advancement.
660    ///
661    /// The watermark is the latest sequence (`seq`) this query round safely read
662    /// for each participating region. Flow uses it to decide where the next
663    /// incremental round can resume.
664    ///
665    /// - `Some(seq)`: the query proved it safely read up to `seq`; downstream
666    ///   may advance the checkpoint to this value.
667    /// - `None`: the region participated but the query could not prove a safe
668    ///   read upper-bound, so the checkpoint must not advance for this region.
669    ///
670    /// Omitted when empty for backward compatibility.
671    #[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
682/// Determines if a metric name represents a time measurement that should be formatted.
683fn is_time_metric(metric_name: &str) -> bool {
684    metric_name.contains("elapsed") || metric_name.contains("time") || metric_name.contains("cost")
685}
686
687/// Determines if a metric name represents a bytes measurement that should be formatted.
688fn 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
696/// Only display `plan_metrics` with indent `  ` (2 spaces).
697impl 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    /// The plan name
731    pub plan: String,
732    /// The stable execution plan name.
733    #[serde(default)]
734    pub plan_name: String,
735    /// The level of the plan, starts from 0
736    pub level: usize,
737    /// An ordered key-value list of metrics.
738    /// Key is metric label and value is metric value.
739    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    // This is not supported for lazy stream.
803    #[inline]
804    fn size_hint(&self) -> (usize, Option<usize>) {
805        (0, None)
806    }
807}
808
809/// Custom cast function that handles Map -> Binary (JSON) conversion
810fn 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
824/// Convert a Map array to a Binary array containing JSON data
825fn 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            // Extract the map entry at index i
844            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            // Convert to JSON object
853            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), // nulls
1334            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}