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::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/// Casts the `RecordBatch`es of `stream` against the `output_schema`.
63#[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
182/// Greptime SendableRecordBatchStream -> DataFusion RecordBatchStream.
183/// The reverse one is [RecordBatchStreamAdapter].
184pub 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
220/// DataFusion [SendableRecordBatchStream](DfSendableRecordBatchStream) -> Greptime [RecordBatchStream].
221/// The reverse one is [DfRecordBatchStreamAdapter].
222/// It can collect metrics from DataFusion execution plan.
223pub struct RecordBatchStreamAdapter {
224    schema: SchemaRef,
225    stream: DfSendableRecordBatchStream,
226    metrics: Option<BaselineMetrics>,
227    /// Aggregated plan-level metrics. Resolved after an [ExecutionPlan] is finished.
228    metrics_2: Metrics,
229    query_load_region_id: Option<u64>,
230    query_stat_counters: Option<RegionQueryStatCounters>,
231    /// Display plan and metrics in verbose mode.
232    explain_verbose: bool,
233    span: Span,
234}
235
236/// Query statistic counters owned by a region.
237#[derive(Debug, Clone)]
238pub struct RegionQueryStatCounters {
239    /// The total query CPU time in nanoseconds.
240    pub query_cpu_time: Arc<AtomicU64>,
241    /// The total scanned bytes.
242    pub query_scanned_bytes: Arc<AtomicU64>,
243}
244
245/// Json encoded metrics. Contains metric from a whole plan tree.
246enum 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    /// Set the verbose mode for displaying plan and metrics.
318    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            // Verbose in-progress snapshots have the same complete topology as
337            // the final snapshot, but use compact (Default) plan formatting.
338            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
380/// Extracts total `output_bytes` from region scan plan nodes.
381pub 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
400/// Collects the complete plan metrics used by terminal metrics and verbose analyze output.
401fn 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
412/// Collects the minimal metrics needed for query-load reporting before EOF.
413///
414/// This intentionally avoids [`MetricCollector`]'s full per-plan aggregation,
415/// sorting, and plan formatting work on the normal query hot path. The result
416/// is still enough for early-stop/cancellation paths to report per-region CPU
417/// time, scanned bytes, and physical-region attribution.
418fn 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
430/// Recursively walks the physical plan and reads raw metric values without
431/// formatting plan nodes.
432fn 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
555/// An [ExecutionPlanVisitor] to collect metrics from a [ExecutionPlan].
556pub 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        // skip if no metric available
577        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        // scrape plan metrics
589        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            // aggregate high-level metrics
605            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
627/// Returns a single-line summary of the root of the plan.
628/// If the `verbose` flag is set, it will display detailed information about the plan.
629fn 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/// [`RecordBatchMetrics`] carrys metrics value
651/// from datanode to frontend through gRPC
652#[derive(serde::Serialize, serde::Deserialize, Default, Debug, Clone)]
653pub struct RecordBatchMetrics {
654    // High-level aggregated metrics
655    /// CPU consumption in nanoseconds
656    pub elapsed_compute: usize,
657    /// Memory used by the plan in bytes
658    pub memory_usage: usize,
659    // Detailed per-plan metrics
660    /// An ordered list of plan metrics, from top to bottom in post-order.
661    pub plan_metrics: Vec<PlanMetrics>,
662    /// Region id that should receive query-load metrics for this scan.
663    #[serde(default, skip_serializing_if = "Option::is_none")]
664    pub query_load_region_id: Option<u64>,
665    /// Per-region watermark for incremental-read checkpoint advancement.
666    ///
667    /// The watermark is the latest sequence (`seq`) this query round safely read
668    /// for each participating region. Flow uses it to decide where the next
669    /// incremental round can resume.
670    ///
671    /// - `Some(seq)`: the query proved it safely read up to `seq`; downstream
672    ///   may advance the checkpoint to this value.
673    /// - `None`: the region participated but the query could not prove a safe
674    ///   read upper-bound, so the checkpoint must not advance for this region.
675    ///
676    /// Omitted when empty for backward compatibility.
677    #[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
688/// Determines if a metric name represents a time measurement that should be formatted.
689fn is_time_metric(metric_name: &str) -> bool {
690    metric_name.contains("elapsed") || metric_name.contains("time") || metric_name.contains("cost")
691}
692
693/// Determines if a metric name represents a bytes measurement that should be formatted.
694fn 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
702/// Only display `plan_metrics` with indent `  ` (2 spaces).
703impl 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    /// The plan name
737    pub plan: String,
738    /// The stable execution plan name.
739    #[serde(default)]
740    pub plan_name: String,
741    /// The level of the plan, starts from 0
742    pub level: usize,
743    /// An ordered key-value list of metrics.
744    /// Key is metric label and value is metric value.
745    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    // This is not supported for lazy stream.
809    #[inline]
810    fn size_hint(&self) -> (usize, Option<usize>) {
811        (0, None)
812    }
813}
814
815/// Custom cast function that handles Map -> Binary (JSON) conversion
816fn 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
830/// Convert a Map array to a Binary array containing JSON data
831fn 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            // Extract the map entry at index i
850            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            // Convert to JSON object
859            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), // nulls
1345            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}