Skip to main content

mito2/read/
scan_util.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
15//! Utilities for scanners.
16
17use std::collections::{BinaryHeap, HashMap, VecDeque};
18use std::fmt;
19use std::pin::Pin;
20use std::sync::{Arc, Mutex};
21use std::task::{Context, Poll};
22use std::time::{Duration, Instant};
23
24use async_stream::try_stream;
25use common_telemetry::tracing;
26use datafusion::physical_plan::metrics::{ExecutionPlanMetricsSet, MetricBuilder, Time};
27use datatypes::arrow::record_batch::RecordBatch;
28use datatypes::timestamp::timestamp_array_to_primitive;
29use futures::Stream;
30use prometheus::IntGauge;
31use smallvec::SmallVec;
32use snafu::ResultExt;
33use store_api::storage::{RegionId, SequenceRange, TimeSeriesRowSelector};
34
35use crate::error::{ComputeArrowSnafu, Result};
36use crate::memtable::MemScanMetrics;
37use crate::metrics::{
38    IN_PROGRESS_SCAN, PRECISE_FILTER_ROWS_TOTAL, READ_BATCHES_RETURN, READ_ROW_GROUPS_TOTAL,
39    READ_ROWS_IN_ROW_GROUP_TOTAL, READ_ROWS_RETURN, READ_STAGE_ELAPSED,
40};
41use crate::read::dedup::{DedupMetrics, DedupMetricsReport};
42use crate::read::flat_merge::{MergeMetrics, MergeMetricsReport};
43use crate::read::pruner::PartitionPruner;
44use crate::read::range::{RangeMeta, RowGroupIndex};
45use crate::read::scan_region::StreamContext;
46use crate::read::{BoxedRecordBatchStream, ScannerMetrics};
47use crate::sst::file::{FileTimeRange, RegionFileId};
48use crate::sst::index::bloom_filter::applier::BloomFilterIndexApplyMetrics;
49use crate::sst::index::fulltext_index::applier::FulltextIndexApplyMetrics;
50use crate::sst::index::inverted_index::applier::InvertedIndexApplyMetrics;
51use crate::sst::parquet::file_range::{FileRange, PreFilterMode};
52use crate::sst::parquet::flat_format::{sequence_column_index, time_index_column_index};
53use crate::sst::parquet::reader::{MetadataCacheMetrics, ReaderFilterMetrics, ReaderMetrics};
54use crate::sst::parquet::row_group::ParquetFetchMetrics;
55use crate::sst::parquet::{DEFAULT_READ_BATCH_SIZE, DEFAULT_ROW_GROUP_SIZE};
56
57/// Per-file scan metrics.
58#[derive(Default, Clone)]
59pub struct FileScanMetrics {
60    /// Number of ranges (row groups) read from this file.
61    pub num_ranges: usize,
62    /// Number of rows read from this file.
63    pub num_rows: usize,
64    /// Time spent building file ranges/parts (file-level preparation).
65    pub build_part_cost: Duration,
66    /// Time spent building readers for this file (accumulated across all ranges).
67    pub build_reader_cost: Duration,
68    /// Time spent scanning this file (accumulated across all ranges).
69    pub scan_cost: Duration,
70}
71
72impl fmt::Debug for FileScanMetrics {
73    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
74        write!(f, "{{\"build_part_cost\":\"{:?}\"", self.build_part_cost)?;
75
76        if self.num_ranges > 0 {
77            write!(f, ", \"num_ranges\":{}", self.num_ranges)?;
78        }
79        if self.num_rows > 0 {
80            write!(f, ", \"num_rows\":{}", self.num_rows)?;
81        }
82        if !self.build_reader_cost.is_zero() {
83            write!(
84                f,
85                ", \"build_reader_cost\":\"{:?}\"",
86                self.build_reader_cost
87            )?;
88        }
89        if !self.scan_cost.is_zero() {
90            write!(f, ", \"scan_cost\":\"{:?}\"", self.scan_cost)?;
91        }
92
93        write!(f, "}}")
94    }
95}
96
97impl FileScanMetrics {
98    /// Merges another FileMetrics into this one.
99    pub(crate) fn merge_from(&mut self, other: &FileScanMetrics) {
100        self.num_ranges += other.num_ranges;
101        self.num_rows += other.num_rows;
102        self.build_part_cost += other.build_part_cost;
103        self.build_reader_cost += other.build_reader_cost;
104        self.scan_cost += other.scan_cost;
105    }
106}
107
108/// Verbose scan metrics for a partition.
109#[derive(Default)]
110pub(crate) struct ScanMetricsSet {
111    /// Duration to prepare the scan task.
112    prepare_scan_cost: Duration,
113    /// Duration to build the (merge) reader.
114    build_reader_cost: Duration,
115    /// Duration to scan data.
116    scan_cost: Duration,
117    /// Duration while waiting for `yield`.
118    yield_cost: Duration,
119    /// Duration to convert [`Batch`]es.
120    convert_cost: Option<Time>,
121    /// Duration of the scan.
122    total_cost: Duration,
123    /// Number of rows returned.
124    num_rows: usize,
125    /// Number of batches returned.
126    num_batches: usize,
127    /// Number of mem ranges scanned.
128    num_mem_ranges: usize,
129    /// Number of file ranges scanned.
130    num_file_ranges: usize,
131
132    // Memtable related metrics:
133    /// Duration to scan memtables.
134    mem_scan_cost: Duration,
135    /// Number of rows read from memtables.
136    mem_rows: usize,
137    /// Number of batches read from memtables.
138    mem_batches: usize,
139    /// Number of series read from memtables.
140    mem_series: usize,
141    /// Duration of prefilter in memtable scan.
142    mem_prefilter_cost: Duration,
143    /// Number of rows filtered by prefilter in memtable scan.
144    mem_prefilter_rows_filtered: usize,
145
146    // SST related metrics:
147    /// Duration to build file ranges.
148    build_parts_cost: Duration,
149    /// Duration to scan SST files.
150    sst_scan_cost: Duration,
151    /// Number of row groups before filtering.
152    rg_total: usize,
153    /// Number of row groups filtered by fulltext index.
154    rg_fulltext_filtered: usize,
155    /// Number of row groups filtered by inverted index.
156    rg_inverted_filtered: usize,
157    /// Number of row groups filtered by min-max index.
158    rg_minmax_filtered: usize,
159    /// Number of row groups filtered by bloom filter index.
160    rg_bloom_filtered: usize,
161    /// Number of rows in row group before filtering.
162    rows_before_filter: usize,
163    /// Number of rows in row group filtered by fulltext index.
164    rows_fulltext_filtered: usize,
165    /// Number of rows in row group filtered by inverted index.
166    rows_inverted_filtered: usize,
167    /// Number of rows in row group filtered by bloom filter index.
168    rows_bloom_filtered: usize,
169    /// Number of rows filtered by precise filter.
170    rows_precise_filtered: usize,
171    /// Number of index result cache hits for fulltext index.
172    fulltext_index_cache_hit: usize,
173    /// Number of index result cache misses for fulltext index.
174    fulltext_index_cache_miss: usize,
175    /// Number of index result cache hits for inverted index.
176    inverted_index_cache_hit: usize,
177    /// Number of index result cache misses for inverted index.
178    inverted_index_cache_miss: usize,
179    /// Number of index result cache hits for bloom filter index.
180    bloom_filter_cache_hit: usize,
181    /// Number of index result cache misses for bloom filter index.
182    bloom_filter_cache_miss: usize,
183    /// Number of index result cache hits for minmax pruning.
184    minmax_cache_hit: usize,
185    /// Number of index result cache misses for minmax pruning.
186    minmax_cache_miss: usize,
187    /// Number of pruner builder cache hits.
188    pruner_cache_hit: usize,
189    /// Number of pruner builder cache misses.
190    pruner_cache_miss: usize,
191    /// Duration spent waiting for pruner to build file ranges.
192    pruner_prune_cost: Duration,
193    /// Number of files filtered by manifest time-range pruning.
194    files_time_range_pruned: usize,
195    /// Number of record batches read from SST.
196    num_sst_record_batches: usize,
197    /// Number of batches decoded from SST.
198    num_sst_batches: usize,
199    /// Number of rows read from SST.
200    num_sst_rows: usize,
201
202    /// Elapsed time before the first poll operation.
203    first_poll: Duration,
204
205    /// Number of send timeout in SeriesScan.
206    num_series_send_timeout: usize,
207    /// Number of send full in SeriesScan.
208    num_series_send_full: usize,
209    /// Number of rows the series distributor scanned.
210    num_distributor_rows: usize,
211    /// Number of batches the series distributor scanned.
212    num_distributor_batches: usize,
213    /// Duration of the series distributor to scan.
214    distributor_scan_cost: Duration,
215    /// Duration of the series distributor to yield.
216    distributor_yield_cost: Duration,
217    /// Duration spent in divider operations.
218    distributor_divider_cost: Duration,
219
220    /// Merge metrics.
221    merge_metrics: MergeMetrics,
222    /// Dedup metrics.
223    dedup_metrics: DedupMetrics,
224
225    /// The stream reached EOF
226    stream_eof: bool,
227
228    // Optional verbose metrics:
229    /// Inverted index apply metrics.
230    inverted_index_apply_metrics: Option<InvertedIndexApplyMetrics>,
231    /// Bloom filter index apply metrics.
232    bloom_filter_apply_metrics: Option<BloomFilterIndexApplyMetrics>,
233    /// Fulltext index apply metrics.
234    fulltext_index_apply_metrics: Option<FulltextIndexApplyMetrics>,
235    /// Parquet fetch metrics.
236    fetch_metrics: Option<ParquetFetchMetrics>,
237    /// Metadata cache metrics.
238    metadata_cache_metrics: Option<MetadataCacheMetrics>,
239    /// Per-file scan metrics, only populated when explain_verbose is true.
240    per_file_metrics: Option<HashMap<RegionFileId, FileScanMetrics>>,
241
242    /// Current memory usage for file range builders.
243    build_ranges_mem_size: isize,
244    /// Peak memory usage for file range builders.
245    build_ranges_peak_mem_size: isize,
246    /// Current number of file range builders.
247    num_range_builders: isize,
248    /// Peak number of file range builders.
249    num_peak_range_builders: isize,
250    /// Total bytes added to the range cache during this scan.
251    range_cache_size: usize,
252    /// Number of range cache hits during this scan.
253    range_cache_hit: usize,
254    /// Number of range cache misses during this scan.
255    range_cache_miss: usize,
256}
257
258/// Wrapper for file metrics that compares by total cost in reverse order.
259/// This allows using BinaryHeap as a min-heap for efficient top-K selection.
260struct CompareCostReverse<'a> {
261    total_cost: Duration,
262    file_id: RegionFileId,
263    metrics: &'a FileScanMetrics,
264}
265
266impl Ord for CompareCostReverse<'_> {
267    fn cmp(&self, other: &Self) -> std::cmp::Ordering {
268        // Reverse comparison: smaller costs are "greater"
269        other.total_cost.cmp(&self.total_cost)
270    }
271}
272
273impl PartialOrd for CompareCostReverse<'_> {
274    fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
275        Some(self.cmp(other))
276    }
277}
278
279impl Eq for CompareCostReverse<'_> {}
280
281impl PartialEq for CompareCostReverse<'_> {
282    fn eq(&self, other: &Self) -> bool {
283        self.total_cost == other.total_cost
284    }
285}
286
287impl fmt::Debug for ScanMetricsSet {
288    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
289        let ScanMetricsSet {
290            prepare_scan_cost,
291            build_reader_cost,
292            scan_cost,
293            yield_cost,
294            convert_cost,
295            total_cost,
296            num_rows,
297            num_batches,
298            num_mem_ranges,
299            num_file_ranges,
300            build_parts_cost,
301            sst_scan_cost,
302            rg_total,
303            rg_fulltext_filtered,
304            rg_inverted_filtered,
305            rg_minmax_filtered,
306            rg_bloom_filtered,
307            rows_before_filter,
308            rows_fulltext_filtered,
309            rows_inverted_filtered,
310            rows_bloom_filtered,
311            rows_precise_filtered,
312            fulltext_index_cache_hit,
313            fulltext_index_cache_miss,
314            inverted_index_cache_hit,
315            inverted_index_cache_miss,
316            bloom_filter_cache_hit,
317            bloom_filter_cache_miss,
318            minmax_cache_hit,
319            minmax_cache_miss,
320            pruner_cache_hit,
321            pruner_cache_miss,
322            pruner_prune_cost,
323            files_time_range_pruned,
324            num_sst_record_batches,
325            num_sst_batches,
326            num_sst_rows,
327            first_poll,
328            num_series_send_timeout,
329            num_series_send_full,
330            num_distributor_rows,
331            num_distributor_batches,
332            distributor_scan_cost,
333            distributor_yield_cost,
334            distributor_divider_cost,
335            merge_metrics,
336            dedup_metrics,
337            stream_eof,
338            mem_scan_cost,
339            mem_rows,
340            mem_batches,
341            mem_series,
342            mem_prefilter_cost,
343            mem_prefilter_rows_filtered,
344            inverted_index_apply_metrics,
345            bloom_filter_apply_metrics,
346            fulltext_index_apply_metrics,
347            fetch_metrics,
348            metadata_cache_metrics,
349            per_file_metrics,
350            build_ranges_mem_size: _,
351            build_ranges_peak_mem_size,
352            num_range_builders: _,
353            num_peak_range_builders,
354            range_cache_size,
355            range_cache_hit,
356            range_cache_miss,
357        } = self;
358
359        // Write core metrics
360        write!(
361            f,
362            "{{\"prepare_scan_cost\":\"{prepare_scan_cost:?}\", \
363            \"build_reader_cost\":\"{build_reader_cost:?}\", \
364            \"scan_cost\":\"{scan_cost:?}\", \
365            \"yield_cost\":\"{yield_cost:?}\", \
366            \"total_cost\":\"{total_cost:?}\", \
367            \"num_rows\":{num_rows}, \
368            \"num_batches\":{num_batches}, \
369            \"num_mem_ranges\":{num_mem_ranges}, \
370            \"num_file_ranges\":{num_file_ranges}, \
371            \"build_parts_cost\":\"{build_parts_cost:?}\", \
372            \"sst_scan_cost\":\"{sst_scan_cost:?}\", \
373            \"rg_total\":{rg_total}, \
374            \"rows_before_filter\":{rows_before_filter}, \
375            \"num_sst_record_batches\":{num_sst_record_batches}, \
376            \"num_sst_batches\":{num_sst_batches}, \
377            \"num_sst_rows\":{num_sst_rows}, \
378            \"first_poll\":\"{first_poll:?}\""
379        )?;
380
381        // Write convert_cost if present
382        if let Some(time) = convert_cost {
383            let duration = Duration::from_nanos(time.value() as u64);
384            write!(f, ", \"convert_cost\":\"{duration:?}\"")?;
385        }
386
387        // Write non-zero filter counters
388        if *files_time_range_pruned > 0 {
389            write!(f, ", \"files_time_range_pruned\":{files_time_range_pruned}")?;
390        }
391        if *rg_fulltext_filtered > 0 {
392            write!(f, ", \"rg_fulltext_filtered\":{rg_fulltext_filtered}")?;
393        }
394        if *rg_inverted_filtered > 0 {
395            write!(f, ", \"rg_inverted_filtered\":{rg_inverted_filtered}")?;
396        }
397        if *rg_minmax_filtered > 0 {
398            write!(f, ", \"rg_minmax_filtered\":{rg_minmax_filtered}")?;
399        }
400        if *rg_bloom_filtered > 0 {
401            write!(f, ", \"rg_bloom_filtered\":{rg_bloom_filtered}")?;
402        }
403        if *rows_fulltext_filtered > 0 {
404            write!(f, ", \"rows_fulltext_filtered\":{rows_fulltext_filtered}")?;
405        }
406        if *rows_inverted_filtered > 0 {
407            write!(f, ", \"rows_inverted_filtered\":{rows_inverted_filtered}")?;
408        }
409        if *rows_bloom_filtered > 0 {
410            write!(f, ", \"rows_bloom_filtered\":{rows_bloom_filtered}")?;
411        }
412        if *rows_precise_filtered > 0 {
413            write!(f, ", \"rows_precise_filtered\":{rows_precise_filtered}")?;
414        }
415        if *fulltext_index_cache_hit > 0 {
416            write!(
417                f,
418                ", \"fulltext_index_cache_hit\":{fulltext_index_cache_hit}"
419            )?;
420        }
421        if *fulltext_index_cache_miss > 0 {
422            write!(
423                f,
424                ", \"fulltext_index_cache_miss\":{fulltext_index_cache_miss}"
425            )?;
426        }
427        if *inverted_index_cache_hit > 0 {
428            write!(
429                f,
430                ", \"inverted_index_cache_hit\":{inverted_index_cache_hit}"
431            )?;
432        }
433        if *inverted_index_cache_miss > 0 {
434            write!(
435                f,
436                ", \"inverted_index_cache_miss\":{inverted_index_cache_miss}"
437            )?;
438        }
439        if *bloom_filter_cache_hit > 0 {
440            write!(f, ", \"bloom_filter_cache_hit\":{bloom_filter_cache_hit}")?;
441        }
442        if *bloom_filter_cache_miss > 0 {
443            write!(f, ", \"bloom_filter_cache_miss\":{bloom_filter_cache_miss}")?;
444        }
445        if *minmax_cache_hit > 0 {
446            write!(f, ", \"minmax_cache_hit\":{minmax_cache_hit}")?;
447        }
448        if *minmax_cache_miss > 0 {
449            write!(f, ", \"minmax_cache_miss\":{minmax_cache_miss}")?;
450        }
451        if *pruner_cache_hit > 0 {
452            write!(f, ", \"pruner_cache_hit\":{pruner_cache_hit}")?;
453        }
454        if *pruner_cache_miss > 0 {
455            write!(f, ", \"pruner_cache_miss\":{pruner_cache_miss}")?;
456        }
457        if !pruner_prune_cost.is_zero() {
458            write!(f, ", \"pruner_prune_cost\":\"{pruner_prune_cost:?}\"")?;
459        }
460
461        // Write non-zero distributor metrics
462        if *num_series_send_timeout > 0 {
463            write!(f, ", \"num_series_send_timeout\":{num_series_send_timeout}")?;
464        }
465        if *num_series_send_full > 0 {
466            write!(f, ", \"num_series_send_full\":{num_series_send_full}")?;
467        }
468        if *num_distributor_rows > 0 {
469            write!(f, ", \"num_distributor_rows\":{num_distributor_rows}")?;
470        }
471        if *num_distributor_batches > 0 {
472            write!(f, ", \"num_distributor_batches\":{num_distributor_batches}")?;
473        }
474        if !distributor_scan_cost.is_zero() {
475            write!(
476                f,
477                ", \"distributor_scan_cost\":\"{distributor_scan_cost:?}\""
478            )?;
479        }
480        if !distributor_yield_cost.is_zero() {
481            write!(
482                f,
483                ", \"distributor_yield_cost\":\"{distributor_yield_cost:?}\""
484            )?;
485        }
486        if !distributor_divider_cost.is_zero() {
487            write!(
488                f,
489                ", \"distributor_divider_cost\":\"{distributor_divider_cost:?}\""
490            )?;
491        }
492
493        // Write non-zero memtable metrics
494        if *mem_rows > 0 {
495            write!(f, ", \"mem_rows\":{mem_rows}")?;
496        }
497        if *mem_batches > 0 {
498            write!(f, ", \"mem_batches\":{mem_batches}")?;
499        }
500        if *mem_series > 0 {
501            write!(f, ", \"mem_series\":{mem_series}")?;
502        }
503        if !mem_scan_cost.is_zero() {
504            write!(f, ", \"mem_scan_cost\":\"{mem_scan_cost:?}\"")?;
505        }
506        if !mem_prefilter_cost.is_zero() {
507            write!(f, ", \"mem_prefilter_cost\":\"{mem_prefilter_cost:?}\"")?;
508        }
509        if *mem_prefilter_rows_filtered > 0 {
510            write!(
511                f,
512                ", \"mem_prefilter_rows_filtered\":{mem_prefilter_rows_filtered}"
513            )?;
514        }
515
516        // Write optional verbose metrics if they are not empty
517        if let Some(metrics) = inverted_index_apply_metrics
518            && !metrics.is_empty()
519        {
520            write!(f, ", \"inverted_index_apply_metrics\":{:?}", metrics)?;
521        }
522        if let Some(metrics) = bloom_filter_apply_metrics
523            && !metrics.is_empty()
524        {
525            write!(f, ", \"bloom_filter_apply_metrics\":{:?}", metrics)?;
526        }
527        if let Some(metrics) = fulltext_index_apply_metrics
528            && !metrics.is_empty()
529        {
530            write!(f, ", \"fulltext_index_apply_metrics\":{:?}", metrics)?;
531        }
532        if let Some(metrics) = fetch_metrics
533            && !metrics.is_empty()
534        {
535            write!(f, ", \"fetch_metrics\":{:?}", metrics)?;
536        }
537        if let Some(metrics) = metadata_cache_metrics
538            && !metrics.is_empty()
539        {
540            write!(f, ", \"metadata_cache_metrics\":{:?}", metrics)?;
541        }
542
543        // Write merge metrics if not empty
544        if !merge_metrics.scan_cost.is_zero() {
545            write!(f, ", \"merge_metrics\":{:?}", merge_metrics)?;
546        }
547
548        // Write dedup metrics if not empty
549        if !dedup_metrics.dedup_cost.is_zero() {
550            write!(f, ", \"dedup_metrics\":{:?}", dedup_metrics)?;
551        }
552
553        // Write top file metrics if present and non-empty
554        if let Some(file_metrics) = per_file_metrics
555            && !file_metrics.is_empty()
556        {
557            // Use min-heap (BinaryHeap with reverse comparison) to keep only top 10
558            let mut heap = BinaryHeap::new();
559            for (file_id, metrics) in file_metrics.iter() {
560                let total_cost =
561                    metrics.build_part_cost + metrics.build_reader_cost + metrics.scan_cost;
562
563                // If the file has been pruned by a pruner, the build part cost may be zero.
564                // If we didn't read any ranges from it, we don't output the file.
565                if total_cost.is_zero() && metrics.num_ranges == 0 {
566                    continue;
567                }
568
569                if heap.len() < 10 {
570                    // Haven't reached 10 yet, just push
571                    heap.push(CompareCostReverse {
572                        total_cost,
573                        file_id: *file_id,
574                        metrics,
575                    });
576                } else if let Some(min_entry) = heap.peek() {
577                    // If current cost is higher than the minimum in our top-10, replace it
578                    if total_cost > min_entry.total_cost {
579                        heap.pop();
580                        heap.push(CompareCostReverse {
581                            total_cost,
582                            file_id: *file_id,
583                            metrics,
584                        });
585                    }
586                }
587            }
588
589            let top_files = heap.into_sorted_vec();
590            write!(f, ", \"top_file_metrics\": {{")?;
591            for (i, item) in top_files.iter().enumerate() {
592                let CompareCostReverse {
593                    total_cost: _,
594                    file_id,
595                    metrics,
596                } = item;
597                if i > 0 {
598                    write!(f, ", ")?;
599                }
600                write!(f, "\"{}\": {:?}", file_id, metrics)?;
601            }
602            write!(f, "}}")?;
603        }
604
605        if *range_cache_size > 0 {
606            write!(f, ", \"range_cache_size\":{range_cache_size}")?;
607        }
608        if *range_cache_hit > 0 {
609            write!(f, ", \"range_cache_hit\":{range_cache_hit}")?;
610        }
611        if *range_cache_miss > 0 {
612            write!(f, ", \"range_cache_miss\":{range_cache_miss}")?;
613        }
614
615        write!(
616            f,
617            ", \"build_ranges_peak_mem_size\":{build_ranges_peak_mem_size}, \
618             \"num_peak_range_builders\":{num_peak_range_builders}, \
619             \"stream_eof\":{stream_eof}}}"
620        )
621    }
622}
623impl ScanMetricsSet {
624    /// Attaches the `prepare_scan_cost` to the metrics set.
625    fn with_prepare_scan_cost(mut self, cost: Duration) -> Self {
626        self.prepare_scan_cost += cost;
627        self
628    }
629
630    /// Attaches the `convert_cost` to the metrics set.
631    fn with_convert_cost(mut self, time: Time) -> Self {
632        self.convert_cost = Some(time);
633        self
634    }
635
636    /// Merges the local scanner metrics.
637    fn merge_scanner_metrics(&mut self, other: &ScannerMetrics) {
638        let ScannerMetrics {
639            scan_cost,
640            yield_cost,
641            num_batches,
642            num_rows,
643        } = other;
644
645        self.scan_cost += *scan_cost;
646        self.yield_cost += *yield_cost;
647        self.num_rows += *num_rows;
648        self.num_batches += *num_batches;
649    }
650
651    /// Merges the local reader metrics.
652    fn merge_reader_metrics(&mut self, other: &ReaderMetrics) {
653        let ReaderMetrics {
654            build_cost,
655            filter_metrics:
656                ReaderFilterMetrics {
657                    rg_total,
658                    rg_fulltext_filtered,
659                    rg_inverted_filtered,
660                    rg_minmax_filtered,
661                    rg_bloom_filtered,
662                    rows_total,
663                    rows_fulltext_filtered,
664                    rows_inverted_filtered,
665                    rows_bloom_filtered,
666                    rows_precise_filtered,
667                    fulltext_index_cache_hit,
668                    fulltext_index_cache_miss,
669                    inverted_index_cache_hit,
670                    inverted_index_cache_miss,
671                    bloom_filter_cache_hit,
672                    bloom_filter_cache_miss,
673                    minmax_cache_hit,
674                    minmax_cache_miss,
675                    pruner_cache_hit,
676                    pruner_cache_miss,
677                    pruner_prune_cost,
678                    files_time_range_pruned,
679                    inverted_index_apply_metrics,
680                    bloom_filter_apply_metrics,
681                    fulltext_index_apply_metrics,
682                },
683            num_record_batches,
684            num_batches,
685            num_rows,
686            scan_cost,
687            metadata_cache_metrics,
688            fetch_metrics,
689            metadata_mem_size,
690            num_range_builders,
691        } = other;
692
693        self.build_parts_cost += *build_cost;
694        self.sst_scan_cost += *scan_cost;
695
696        self.files_time_range_pruned += *files_time_range_pruned;
697
698        self.rg_total += *rg_total;
699        self.rg_fulltext_filtered += *rg_fulltext_filtered;
700        self.rg_inverted_filtered += *rg_inverted_filtered;
701        self.rg_minmax_filtered += *rg_minmax_filtered;
702        self.rg_bloom_filtered += *rg_bloom_filtered;
703
704        self.rows_before_filter += *rows_total;
705        self.rows_fulltext_filtered += *rows_fulltext_filtered;
706        self.rows_inverted_filtered += *rows_inverted_filtered;
707        self.rows_bloom_filtered += *rows_bloom_filtered;
708        self.rows_precise_filtered += *rows_precise_filtered;
709
710        self.fulltext_index_cache_hit += *fulltext_index_cache_hit;
711        self.fulltext_index_cache_miss += *fulltext_index_cache_miss;
712        self.inverted_index_cache_hit += *inverted_index_cache_hit;
713        self.inverted_index_cache_miss += *inverted_index_cache_miss;
714        self.bloom_filter_cache_hit += *bloom_filter_cache_hit;
715        self.bloom_filter_cache_miss += *bloom_filter_cache_miss;
716        self.minmax_cache_hit += *minmax_cache_hit;
717        self.minmax_cache_miss += *minmax_cache_miss;
718        self.pruner_cache_hit += *pruner_cache_hit;
719        self.pruner_cache_miss += *pruner_cache_miss;
720        self.pruner_prune_cost += *pruner_prune_cost;
721
722        self.num_sst_record_batches += *num_record_batches;
723        self.num_sst_batches += *num_batches;
724        self.num_sst_rows += *num_rows;
725
726        // Merge optional verbose metrics
727        if let Some(metrics) = inverted_index_apply_metrics {
728            self.inverted_index_apply_metrics
729                .get_or_insert_with(InvertedIndexApplyMetrics::default)
730                .merge_from(metrics);
731        }
732        if let Some(metrics) = bloom_filter_apply_metrics {
733            self.bloom_filter_apply_metrics
734                .get_or_insert_with(BloomFilterIndexApplyMetrics::default)
735                .merge_from(metrics);
736        }
737        if let Some(metrics) = fulltext_index_apply_metrics {
738            self.fulltext_index_apply_metrics
739                .get_or_insert_with(FulltextIndexApplyMetrics::default)
740                .merge_from(metrics);
741        }
742        if let Some(metrics) = fetch_metrics {
743            self.fetch_metrics
744                .get_or_insert_with(ParquetFetchMetrics::default)
745                .merge_from(metrics);
746        }
747        self.metadata_cache_metrics
748            .get_or_insert_with(MetadataCacheMetrics::default)
749            .merge_from(metadata_cache_metrics);
750
751        // Track memory usage and update peak.
752        self.build_ranges_mem_size += *metadata_mem_size;
753        if self.build_ranges_mem_size > self.build_ranges_peak_mem_size {
754            self.build_ranges_peak_mem_size = self.build_ranges_mem_size;
755        }
756
757        // Track number of builders and update peak.
758        self.num_range_builders += *num_range_builders;
759        if self.num_range_builders > self.num_peak_range_builders {
760            self.num_peak_range_builders = self.num_range_builders;
761        }
762    }
763
764    /// Merges per-file metrics.
765    fn merge_per_file_metrics(&mut self, other: &HashMap<RegionFileId, FileScanMetrics>) {
766        let self_file_metrics = self.per_file_metrics.get_or_insert_with(HashMap::new);
767        for (file_id, metrics) in other {
768            self_file_metrics
769                .entry(*file_id)
770                .or_default()
771                .merge_from(metrics);
772        }
773    }
774
775    /// Sets distributor metrics.
776    fn set_distributor_metrics(&mut self, distributor_metrics: &SeriesDistributorMetrics) {
777        let SeriesDistributorMetrics {
778            num_series_send_timeout,
779            num_series_send_full,
780            num_rows,
781            num_batches,
782            scan_cost,
783            yield_cost,
784            divider_cost,
785        } = distributor_metrics;
786
787        self.num_series_send_timeout += *num_series_send_timeout;
788        self.num_series_send_full += *num_series_send_full;
789        self.num_distributor_rows += *num_rows;
790        self.num_distributor_batches += *num_batches;
791        self.distributor_scan_cost += *scan_cost;
792        self.distributor_yield_cost += *yield_cost;
793        self.distributor_divider_cost += *divider_cost;
794    }
795
796    /// Observes metrics.
797    fn observe_metrics(&self) {
798        READ_STAGE_ELAPSED
799            .with_label_values(&["prepare_scan"])
800            .observe(self.prepare_scan_cost.as_secs_f64());
801        READ_STAGE_ELAPSED
802            .with_label_values(&["build_reader"])
803            .observe(self.build_reader_cost.as_secs_f64());
804        READ_STAGE_ELAPSED
805            .with_label_values(&["scan"])
806            .observe(self.scan_cost.as_secs_f64());
807        READ_STAGE_ELAPSED
808            .with_label_values(&["yield"])
809            .observe(self.yield_cost.as_secs_f64());
810        if let Some(time) = &self.convert_cost {
811            READ_STAGE_ELAPSED
812                .with_label_values(&["convert"])
813                .observe(Duration::from_nanos(time.value() as u64).as_secs_f64());
814        }
815        READ_STAGE_ELAPSED
816            .with_label_values(&["total"])
817            .observe(self.total_cost.as_secs_f64());
818        READ_ROWS_RETURN.observe(self.num_rows as f64);
819        READ_BATCHES_RETURN.observe(self.num_batches as f64);
820
821        READ_STAGE_ELAPSED
822            .with_label_values(&["build_parts"])
823            .observe(self.build_parts_cost.as_secs_f64());
824
825        READ_ROW_GROUPS_TOTAL
826            .with_label_values(&["before_filtering"])
827            .inc_by(self.rg_total as u64);
828        READ_ROW_GROUPS_TOTAL
829            .with_label_values(&["fulltext_index_filtered"])
830            .inc_by(self.rg_fulltext_filtered as u64);
831        READ_ROW_GROUPS_TOTAL
832            .with_label_values(&["inverted_index_filtered"])
833            .inc_by(self.rg_inverted_filtered as u64);
834        READ_ROW_GROUPS_TOTAL
835            .with_label_values(&["minmax_index_filtered"])
836            .inc_by(self.rg_minmax_filtered as u64);
837        READ_ROW_GROUPS_TOTAL
838            .with_label_values(&["bloom_filter_index_filtered"])
839            .inc_by(self.rg_bloom_filtered as u64);
840
841        PRECISE_FILTER_ROWS_TOTAL
842            .with_label_values(&["parquet"])
843            .inc_by(self.rows_precise_filtered as u64);
844        READ_ROWS_IN_ROW_GROUP_TOTAL
845            .with_label_values(&["before_filtering"])
846            .inc_by(self.rows_before_filter as u64);
847        READ_ROWS_IN_ROW_GROUP_TOTAL
848            .with_label_values(&["fulltext_index_filtered"])
849            .inc_by(self.rows_fulltext_filtered as u64);
850        READ_ROWS_IN_ROW_GROUP_TOTAL
851            .with_label_values(&["inverted_index_filtered"])
852            .inc_by(self.rows_inverted_filtered as u64);
853        READ_ROWS_IN_ROW_GROUP_TOTAL
854            .with_label_values(&["bloom_filter_index_filtered"])
855            .inc_by(self.rows_bloom_filtered as u64);
856    }
857}
858
859struct PartitionMetricsInner {
860    region_id: RegionId,
861    /// Index of the partition to scan.
862    partition: usize,
863    /// Label to distinguish different scan operation.
864    scanner_type: &'static str,
865    /// Query start time.
866    query_start: Instant,
867    /// Whether to use verbose logging.
868    explain_verbose: bool,
869    /// Verbose scan metrics that only log to debug logs by default.
870    metrics: Mutex<ScanMetricsSet>,
871    in_progress_scan: IntGauge,
872
873    // Normal metrics that always report to the [ExecutionPlanMetricsSet]:
874    /// Duration to build file ranges.
875    build_parts_cost: Time,
876    /// Duration to build the (merge) reader.
877    build_reader_cost: Time,
878    /// Duration to scan data.
879    scan_cost: Time,
880    /// Duration while waiting for `yield`.
881    yield_cost: Time,
882    /// Duration to convert [`Batch`]es.
883    convert_cost: Time,
884    /// Aggregated compute time reported to DataFusion.
885    elapsed_compute: Time,
886}
887
888impl PartitionMetricsInner {
889    fn on_finish(&self, stream_eof: bool) {
890        let mut metrics = self.metrics.lock().unwrap();
891        if metrics.total_cost.is_zero() {
892            metrics.total_cost = self.query_start.elapsed();
893        }
894        if !metrics.stream_eof {
895            metrics.stream_eof = stream_eof;
896        }
897    }
898}
899
900impl MergeMetricsReport for PartitionMetricsInner {
901    fn report(&self, metrics: &mut MergeMetrics) {
902        let mut scan_metrics = self.metrics.lock().unwrap();
903        // Merge the metrics into scan_metrics
904        scan_metrics.merge_metrics.merge(metrics);
905
906        // Reset the input metrics
907        *metrics = MergeMetrics::default();
908    }
909}
910
911impl DedupMetricsReport for PartitionMetricsInner {
912    fn report(&self, metrics: &mut DedupMetrics) {
913        let mut scan_metrics = self.metrics.lock().unwrap();
914        // Merge the metrics into scan_metrics
915        scan_metrics.dedup_metrics.merge(metrics);
916
917        // Reset the input metrics
918        *metrics = DedupMetrics::default();
919    }
920}
921
922impl Drop for PartitionMetricsInner {
923    fn drop(&mut self) {
924        self.on_finish(false);
925        let metrics = self.metrics.lock().unwrap();
926        metrics.observe_metrics();
927        self.in_progress_scan.dec();
928
929        if self.explain_verbose {
930            common_telemetry::info!(
931                "{} finished, region_id: {}, partition: {}, scan_metrics: {:?}",
932                self.scanner_type,
933                self.region_id,
934                self.partition,
935                metrics,
936            );
937        } else {
938            common_telemetry::debug!(
939                "{} finished, region_id: {}, partition: {}, scan_metrics: {:?}",
940                self.scanner_type,
941                self.region_id,
942                self.partition,
943                metrics,
944            );
945        }
946    }
947}
948
949/// List of PartitionMetrics.
950#[derive(Default)]
951pub(crate) struct PartitionMetricsList(Mutex<Vec<Option<PartitionMetrics>>>);
952
953impl PartitionMetricsList {
954    /// Sets a new [PartitionMetrics] at the specified partition.
955    pub(crate) fn set(&self, partition: usize, metrics: PartitionMetrics) {
956        let mut list = self.0.lock().unwrap();
957        if list.len() <= partition {
958            list.resize(partition + 1, None);
959        }
960        list[partition] = Some(metrics);
961    }
962
963    /// Format verbose metrics for each partition for explain.
964    pub(crate) fn format_verbose_metrics(&self, f: &mut fmt::Formatter) -> fmt::Result {
965        let list = self.0.lock().unwrap();
966        write!(f, ", \"metrics_per_partition\": ")?;
967        f.debug_list()
968            .entries(list.iter().filter_map(|p| p.as_ref()))
969            .finish()?;
970        write!(f, "}}")
971    }
972}
973
974/// Metrics while reading a partition.
975#[derive(Clone)]
976pub struct PartitionMetrics(Arc<PartitionMetricsInner>);
977
978impl PartitionMetrics {
979    pub(crate) fn new(
980        region_id: RegionId,
981        partition: usize,
982        scanner_type: &'static str,
983        query_start: Instant,
984        explain_verbose: bool,
985        metrics_set: &ExecutionPlanMetricsSet,
986    ) -> Self {
987        let partition_str = partition.to_string();
988        let in_progress_scan = IN_PROGRESS_SCAN.with_label_values(&[scanner_type, &partition_str]);
989        in_progress_scan.inc();
990        let convert_cost = MetricBuilder::new(metrics_set).subset_time("convert_cost", partition);
991        let metrics = ScanMetricsSet::default()
992            .with_prepare_scan_cost(query_start.elapsed())
993            .with_convert_cost(convert_cost.clone());
994        let inner = PartitionMetricsInner {
995            region_id,
996            partition,
997            scanner_type,
998            query_start,
999            explain_verbose,
1000            metrics: Mutex::new(metrics),
1001            in_progress_scan,
1002            build_parts_cost: MetricBuilder::new(metrics_set)
1003                .subset_time("build_parts_cost", partition),
1004            build_reader_cost: MetricBuilder::new(metrics_set)
1005                .subset_time("build_reader_cost", partition),
1006            scan_cost: MetricBuilder::new(metrics_set).subset_time("scan_cost", partition),
1007            yield_cost: MetricBuilder::new(metrics_set).subset_time("yield_cost", partition),
1008            convert_cost,
1009            elapsed_compute: MetricBuilder::new(metrics_set).elapsed_compute(partition),
1010        };
1011        Self(Arc::new(inner))
1012    }
1013
1014    pub(crate) fn on_first_poll(&self) {
1015        let mut metrics = self.0.metrics.lock().unwrap();
1016        metrics.first_poll = self.0.query_start.elapsed();
1017    }
1018
1019    pub(crate) fn inc_num_mem_ranges(&self, num: usize) {
1020        let mut metrics = self.0.metrics.lock().unwrap();
1021        metrics.num_mem_ranges += num;
1022    }
1023
1024    pub fn inc_num_file_ranges(&self, num: usize) {
1025        let mut metrics = self.0.metrics.lock().unwrap();
1026        metrics.num_file_ranges += num;
1027    }
1028
1029    fn record_elapsed_compute(&self, duration: Duration) {
1030        if duration.is_zero() {
1031            return;
1032        }
1033        self.0.elapsed_compute.add_duration(duration);
1034    }
1035
1036    /// Merges `build_reader_cost`.
1037    pub(crate) fn inc_build_reader_cost(&self, cost: Duration) {
1038        self.0.build_reader_cost.add_duration(cost);
1039
1040        let mut metrics = self.0.metrics.lock().unwrap();
1041        metrics.build_reader_cost += cost;
1042    }
1043
1044    pub(crate) fn inc_convert_batch_cost(&self, cost: Duration) {
1045        self.0.convert_cost.add_duration(cost);
1046        self.record_elapsed_compute(cost);
1047    }
1048
1049    /// Reports memtable scan metrics.
1050    pub(crate) fn report_mem_scan_metrics(&self, data: &crate::memtable::MemScanMetricsData) {
1051        let mut metrics = self.0.metrics.lock().unwrap();
1052        metrics.mem_scan_cost += data.scan_cost;
1053        metrics.mem_rows += data.num_rows;
1054        metrics.mem_batches += data.num_batches;
1055        metrics.mem_series += data.total_series;
1056        metrics.mem_prefilter_cost += data.prefilter_cost;
1057        metrics.mem_prefilter_rows_filtered += data.prefilter_rows_filtered;
1058    }
1059
1060    /// Merges [ScannerMetrics], `build_reader_cost`, `scan_cost` and `yield_cost`.
1061    pub(crate) fn merge_metrics(&self, metrics: &ScannerMetrics) {
1062        self.0.scan_cost.add_duration(metrics.scan_cost);
1063        self.record_elapsed_compute(metrics.scan_cost);
1064        self.0.yield_cost.add_duration(metrics.yield_cost);
1065        self.record_elapsed_compute(metrics.yield_cost);
1066
1067        let mut metrics_set = self.0.metrics.lock().unwrap();
1068        metrics_set.merge_scanner_metrics(metrics);
1069    }
1070
1071    /// Merges [ReaderMetrics] and `build_reader_cost`.
1072    pub fn merge_reader_metrics(
1073        &self,
1074        metrics: &ReaderMetrics,
1075        per_file_metrics: Option<&HashMap<RegionFileId, FileScanMetrics>>,
1076    ) {
1077        self.0.build_parts_cost.add_duration(metrics.build_cost);
1078
1079        let mut metrics_set = self.0.metrics.lock().unwrap();
1080        metrics_set.merge_reader_metrics(metrics);
1081
1082        // Merge per-file metrics if provided
1083        if let Some(file_metrics) = per_file_metrics {
1084            metrics_set.merge_per_file_metrics(file_metrics);
1085        }
1086    }
1087
1088    /// Finishes the query.
1089    pub(crate) fn on_finish(&self) {
1090        self.0.on_finish(true);
1091    }
1092
1093    /// Sets the distributor metrics.
1094    pub(crate) fn set_distributor_metrics(&self, metrics: &SeriesDistributorMetrics) {
1095        let mut metrics_set = self.0.metrics.lock().unwrap();
1096        metrics_set.set_distributor_metrics(metrics);
1097    }
1098
1099    /// Returns whether verbose explain is enabled.
1100    pub(crate) fn explain_verbose(&self) -> bool {
1101        self.0.explain_verbose
1102    }
1103
1104    /// Returns a MergeMetricsReport trait object for reporting merge metrics.
1105    pub(crate) fn merge_metrics_reporter(&self) -> Arc<dyn MergeMetricsReport> {
1106        self.0.clone()
1107    }
1108
1109    /// Returns a DedupMetricsReport trait object for reporting dedup metrics.
1110    pub(crate) fn dedup_metrics_reporter(&self) -> Arc<dyn DedupMetricsReport> {
1111        self.0.clone()
1112    }
1113
1114    /// Increments the total bytes added to the range cache.
1115    pub(crate) fn inc_range_cache_size(&self, size: usize) {
1116        let mut metrics = self.0.metrics.lock().unwrap();
1117        metrics.range_cache_size += size;
1118    }
1119
1120    /// Increments the range cache hit counter.
1121    pub(crate) fn inc_range_cache_hit(&self) {
1122        let mut metrics = self.0.metrics.lock().unwrap();
1123        metrics.range_cache_hit += 1;
1124    }
1125
1126    /// Increments the range cache miss counter.
1127    pub(crate) fn inc_range_cache_miss(&self) {
1128        let mut metrics = self.0.metrics.lock().unwrap();
1129        metrics.range_cache_miss += 1;
1130    }
1131}
1132
1133impl fmt::Debug for PartitionMetrics {
1134    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1135        let metrics = self.0.metrics.lock().unwrap();
1136        write!(
1137            f,
1138            r#"{{"partition":{}, "metrics":{:?}}}"#,
1139            self.0.partition, metrics
1140        )
1141    }
1142}
1143
1144/// Metrics for the series distributor.
1145#[derive(Default)]
1146pub(crate) struct SeriesDistributorMetrics {
1147    /// Number of send timeout in SeriesScan.
1148    pub(crate) num_series_send_timeout: usize,
1149    /// Number of send full in SeriesScan.
1150    pub(crate) num_series_send_full: usize,
1151    /// Number of rows the series distributor scanned.
1152    pub(crate) num_rows: usize,
1153    /// Number of batches the series distributor scanned.
1154    pub(crate) num_batches: usize,
1155    /// Duration of the series distributor to scan.
1156    pub(crate) scan_cost: Duration,
1157    /// Duration of the series distributor to yield.
1158    pub(crate) yield_cost: Duration,
1159    /// Duration spent in divider operations.
1160    pub(crate) divider_cost: Duration,
1161}
1162
1163/// Scans memtable ranges at `index` using flat format that returns RecordBatch.
1164#[tracing::instrument(
1165    skip_all,
1166    fields(
1167        region_id = %stream_ctx.input.region_metadata().region_id,
1168        row_group_index = %index.index,
1169        source = "mem_flat"
1170    )
1171)]
1172pub(crate) fn scan_flat_mem_ranges(
1173    stream_ctx: Arc<StreamContext>,
1174    part_metrics: PartitionMetrics,
1175    index: RowGroupIndex,
1176    time_range: FileTimeRange,
1177) -> impl Stream<Item = Result<RecordBatch>> {
1178    try_stream! {
1179        let ranges = stream_ctx.input.build_mem_ranges(index);
1180        part_metrics.inc_num_mem_ranges(ranges.len());
1181        for range in ranges {
1182            let build_reader_start = Instant::now();
1183            let mem_scan_metrics = Some(MemScanMetrics::default());
1184            let mut iter = range.build_record_batch_iter(Some(time_range), mem_scan_metrics.clone())?;
1185            part_metrics.inc_build_reader_cost(build_reader_start.elapsed());
1186
1187            while let Some(record_batch) = iter.next().transpose()? {
1188                yield record_batch;
1189            }
1190
1191            // Report the memtable scan metrics to partition metrics
1192            if let Some(ref metrics) = mem_scan_metrics {
1193                let data = metrics.data();
1194                part_metrics.report_mem_scan_metrics(&data);
1195            }
1196        }
1197    }
1198}
1199
1200/// Files with row count greater than this threshold can contribute to the estimation.
1201const SPLIT_ROW_THRESHOLD: u64 = DEFAULT_ROW_GROUP_SIZE as u64;
1202/// Number of series threshold for splitting batches.
1203const NUM_SERIES_THRESHOLD: u64 = 10240;
1204/// Minimum batch size after splitting. The batch size is less than 60 because a series may only have
1205/// 60 samples per hour.
1206const BATCH_SIZE_THRESHOLD: u64 = 50;
1207
1208/// Returns the estimated rows per batch after splitting if splitting flat record batches
1209/// may improve merge performance. Returns `None` if splitting is not beneficial.
1210pub(crate) fn should_split_flat_batches_for_merge(
1211    stream_ctx: &Arc<StreamContext>,
1212    range_meta: &RangeMeta,
1213) -> Option<usize> {
1214    // Number of files to split and scan.
1215    let mut num_files_to_split = 0;
1216    let mut num_mem_rows = 0;
1217    let mut num_mem_series = 0;
1218    // Total rows and series for estimating batch size after splitting.
1219    let mut total_rows: u64 = 0;
1220    let mut total_series: u64 = 0;
1221    // Checks each file range, returns early if any range is not splittable.
1222    // For mem ranges, we collect the total number of rows and series because the number of rows in a
1223    // mem range may be too small.
1224    for index in &range_meta.row_group_indices {
1225        if stream_ctx.is_mem_range_index(*index) {
1226            let memtable = &stream_ctx.input.memtables[index.index];
1227            // Is mem range
1228            let stats = memtable.stats();
1229            num_mem_rows += stats.num_rows();
1230            num_mem_series += stats.series_count();
1231        } else if stream_ctx.is_file_range_index(*index) {
1232            // This is a file range.
1233            let file_index = index.index - stream_ctx.input.num_memtables();
1234            let file = &stream_ctx.input.files[file_index];
1235            let file_meta = file.meta_ref();
1236            if file_meta.level == 0 {
1237                // Always split level 0 files.
1238                num_files_to_split += 1;
1239                continue;
1240            } else if file_meta.num_rows < SPLIT_ROW_THRESHOLD || file_meta.num_series == 0 {
1241                // If the file doesn't have enough rows, or the number of series is unavailable, skips it.
1242                continue;
1243            }
1244            debug_assert!(file_meta.num_rows > 0);
1245            if !can_split_series(file_meta.num_rows, file_meta.num_series) {
1246                // We can't split batches in a file.
1247                common_telemetry::trace!(
1248                    "Can't split series for file {}, level: {}, num_rows: {}, num_series: {}",
1249                    file_meta.file_id,
1250                    file_meta.level,
1251                    file_meta.num_rows,
1252                    file_meta.num_series,
1253                );
1254                return None;
1255            } else {
1256                num_files_to_split += 1;
1257                total_rows += file.meta_ref().num_rows;
1258                total_series += file.meta_ref().num_series;
1259            }
1260        }
1261        // Skips non-file and non-mem ranges.
1262    }
1263
1264    let should_split = if num_files_to_split > 0 {
1265        // We mainly consider file ranges because they have enough data for sampling.
1266        true
1267    } else if num_mem_series > 0
1268        && num_mem_rows > 0
1269        && can_split_series(num_mem_rows as u64, num_mem_series as u64)
1270    {
1271        total_rows += num_mem_rows as u64;
1272        total_series += num_mem_series as u64;
1273        true
1274    } else {
1275        false
1276    };
1277
1278    if !should_split {
1279        return None;
1280    }
1281
1282    // Estimate rows per batch after splitting.
1283    let estimated_batch_size = if total_series > 0 && total_rows > 0 {
1284        ((total_rows / total_series) as usize).clamp(1, DEFAULT_READ_BATCH_SIZE)
1285    } else {
1286        // No valid estimate available, use a conservative fallback.
1287        DEFAULT_READ_BATCH_SIZE / 4
1288    };
1289    Some(estimated_batch_size)
1290}
1291
1292/// Computes the channel size for parallel scan based on the estimated rows per batch.
1293/// The channel should buffer approximately `2 * DEFAULT_READ_BATCH_SIZE` rows.
1294pub(crate) fn compute_parallel_channel_size(estimated_rows_per_batch: usize) -> usize {
1295    let size = 2 * DEFAULT_READ_BATCH_SIZE / estimated_rows_per_batch.max(1);
1296    size.clamp(2, 64)
1297}
1298
1299/// Computes the average estimated rows per batch across multiple range readers.
1300pub(crate) fn compute_average_batch_size(
1301    estimated_rows_per_batch: impl IntoIterator<Item = usize>,
1302) -> usize {
1303    let mut total = 0usize;
1304    let mut count = 0usize;
1305    for size in estimated_rows_per_batch {
1306        total += size;
1307        count += 1;
1308    }
1309
1310    if count == 0 {
1311        return DEFAULT_READ_BATCH_SIZE;
1312    }
1313
1314    (total / count).clamp(1, DEFAULT_READ_BATCH_SIZE)
1315}
1316
1317fn can_split_series(num_rows: u64, num_series: u64) -> bool {
1318    if num_rows == 0 || num_series == 0 {
1319        return false;
1320    }
1321
1322    // It doesn't have too many series or it will have enough rows for each batch.
1323    num_series < NUM_SERIES_THRESHOLD || num_rows / num_series >= BATCH_SIZE_THRESHOLD
1324}
1325
1326#[cfg(test)]
1327mod split_tests {
1328    use std::sync::Arc;
1329
1330    use common_time::Timestamp;
1331    use smallvec::smallvec;
1332    use store_api::storage::FileId;
1333
1334    use super::*;
1335    use crate::read::flat_projection::FlatProjectionMapper;
1336    use crate::read::range::{RangeMeta, RowGroupIndex, SourceIndex};
1337    use crate::read::scan_region::{ScanInput, StreamContext};
1338    use crate::sst::file::FileHandle;
1339    use crate::test_util::memtable_util::metadata_with_primary_key;
1340    use crate::test_util::scheduler_util::SchedulerEnv;
1341    use crate::test_util::sst_util::sst_file_handle_with_file_id;
1342
1343    async fn new_stream_context_with_files(files: Vec<FileHandle>) -> StreamContext {
1344        let env = SchedulerEnv::new().await;
1345        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
1346        let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
1347        let input = ScanInput::builder(env.access_layer.clone(), mapper)
1348            .with_files(files)
1349            .build();
1350
1351        StreamContext {
1352            input,
1353            ranges: vec![],
1354            query_start: std::time::Instant::now(),
1355        }
1356    }
1357
1358    fn single_file_range_meta() -> RangeMeta {
1359        RangeMeta {
1360            time_range: (
1361                Timestamp::new_millisecond(0),
1362                Timestamp::new_millisecond(1000),
1363            ),
1364            indices: smallvec![SourceIndex {
1365                index: 0,
1366                num_row_groups: 1,
1367            }],
1368            row_group_indices: smallvec![RowGroupIndex {
1369                index: 0,
1370                row_group_index: 0,
1371            }],
1372            num_rows: 1024,
1373        }
1374    }
1375
1376    #[tokio::test]
1377    async fn should_split_level_zero_file_even_when_series_stats_are_missing() {
1378        let mut file = sst_file_handle_with_file_id(FileId::random(), 0, 1000)
1379            .meta_ref()
1380            .clone();
1381        file.level = 0;
1382        file.num_rows = DEFAULT_ROW_GROUP_SIZE as u64;
1383        file.num_row_groups = 1;
1384        file.num_series = 0;
1385
1386        let file = FileHandle::new(file, crate::test_util::new_noop_file_purger());
1387        let stream_ctx = Arc::new(new_stream_context_with_files(vec![file]).await);
1388
1389        assert!(
1390            should_split_flat_batches_for_merge(&stream_ctx, &single_file_range_meta()).is_some()
1391        );
1392    }
1393
1394    #[test]
1395    fn can_split_series_returns_false_for_zero_inputs() {
1396        assert!(!can_split_series(0, 1));
1397        assert!(!can_split_series(1, 0));
1398        assert!(!can_split_series(0, 0));
1399    }
1400}
1401
1402/// Creates a new [ReaderFilterMetrics] with optional apply metrics initialized
1403/// based on the `explain_verbose` flag.
1404pub(crate) fn new_filter_metrics(explain_verbose: bool) -> ReaderFilterMetrics {
1405    if explain_verbose {
1406        ReaderFilterMetrics {
1407            inverted_index_apply_metrics: Some(InvertedIndexApplyMetrics::default()),
1408            bloom_filter_apply_metrics: Some(BloomFilterIndexApplyMetrics::default()),
1409            fulltext_index_apply_metrics: Some(FulltextIndexApplyMetrics::default()),
1410            ..Default::default()
1411        }
1412    } else {
1413        ReaderFilterMetrics::default()
1414    }
1415}
1416
1417/// Scans file ranges at `index` using flat reader that returns RecordBatch.
1418#[tracing::instrument(
1419    skip_all,
1420    fields(
1421        region_id = %stream_ctx.input.region_metadata().region_id,
1422        row_group_index = %index.index,
1423        source = read_type
1424    )
1425)]
1426pub(crate) async fn scan_flat_file_ranges(
1427    stream_ctx: Arc<StreamContext>,
1428    part_metrics: PartitionMetrics,
1429    index: RowGroupIndex,
1430    read_type: &'static str,
1431    partition_pruner: Arc<PartitionPruner>,
1432) -> Result<impl Stream<Item = Result<RecordBatch>>> {
1433    let mut reader_metrics = ReaderMetrics {
1434        filter_metrics: new_filter_metrics(part_metrics.explain_verbose()),
1435        ..Default::default()
1436    };
1437    let ranges = partition_pruner
1438        .build_file_ranges(index, &part_metrics, &mut reader_metrics)
1439        .await?;
1440    part_metrics.inc_num_file_ranges(ranges.len());
1441    part_metrics.merge_reader_metrics(&reader_metrics, None);
1442
1443    // Creates initial per-file metrics with build_part_cost.
1444    let init_per_file_metrics = if part_metrics.explain_verbose() {
1445        let file = stream_ctx.input.file_from_index(index);
1446        let file_id = file.file_id();
1447
1448        let mut map = HashMap::new();
1449        map.insert(
1450            file_id,
1451            FileScanMetrics {
1452                build_part_cost: reader_metrics.build_cost,
1453                ..Default::default()
1454            },
1455        );
1456        Some(map)
1457    } else {
1458        None
1459    };
1460
1461    Ok(build_flat_file_range_scan_stream(
1462        stream_ctx,
1463        part_metrics,
1464        read_type,
1465        ranges,
1466        init_per_file_metrics,
1467    ))
1468}
1469
1470/// Filters a flat-format record batch by the exact sequence range.
1471///
1472/// Returns `None` when the entire batch is filtered out. The sequence column is
1473/// the second-to-last internal column of the flat format
1474/// (`__primary_key`, `__sequence`, `__op_type`), so it must be applied before any
1475/// projection/compat conversion that drops internal columns.
1476///
1477/// `file_sequence_trusted` is the per-file trust decision. Foreign reader
1478/// batches have already been virtualized to the target-local file barrier, so
1479/// they are filtered using the effective batch sequence. Local untrusted files
1480/// pass through because exact capability excludes them before row filtering.
1481pub(crate) fn filter_flat_batch_by_sequence(
1482    record_batch: RecordBatch,
1483    sequence_range: Option<SequenceRange>,
1484    file_sequence_trusted: bool,
1485) -> Result<Option<RecordBatch>> {
1486    let Some(sequence) = sequence_range else {
1487        return Ok(Some(record_batch));
1488    };
1489    if !file_sequence_trusted {
1490        return Ok(Some(record_batch));
1491    }
1492
1493    let num_rows = record_batch.num_rows();
1494    if num_rows == 0 {
1495        return Ok(Some(record_batch));
1496    }
1497    let sequence_column = record_batch.column(sequence_column_index(record_batch.num_columns()));
1498    let predicate = sequence
1499        .filter(sequence_column)
1500        .context(ComputeArrowSnafu)?;
1501    let select_count = predicate.true_count();
1502    if select_count == 0 {
1503        return Ok(None);
1504    }
1505    if select_count == num_rows {
1506        return Ok(Some(record_batch));
1507    }
1508    let filtered_batch = datatypes::arrow::compute::filter_record_batch(&record_batch, &predicate)
1509        .context(ComputeArrowSnafu)?;
1510    Ok(Some(filtered_batch))
1511}
1512
1513/// Build the stream of scanning the input [`FileRange`]s using flat reader that returns RecordBatch.
1514#[tracing::instrument(
1515    skip_all,
1516    fields(read_type = read_type, range_count = ranges.len())
1517)]
1518pub fn build_flat_file_range_scan_stream(
1519    stream_ctx: Arc<StreamContext>,
1520    part_metrics: PartitionMetrics,
1521    read_type: &'static str,
1522    ranges: SmallVec<[FileRange; 2]>,
1523    mut per_file_metrics: Option<HashMap<RegionFileId, FileScanMetrics>>,
1524) -> impl Stream<Item = Result<RecordBatch>> {
1525    try_stream! {
1526        let fetch_metrics = if part_metrics.explain_verbose() {
1527            Some(Arc::new(ParquetFetchMetrics::default()))
1528        } else {
1529            None
1530        };
1531        let reader_metrics = &mut ReaderMetrics {
1532            fetch_metrics: fetch_metrics.clone(),
1533            ..Default::default()
1534        };
1535        for range in ranges {
1536            let build_reader_start = Instant::now();
1537            let Some(mut reader) = range
1538                .flat_reader(
1539                    match stream_ctx.input.series_row_selector {
1540                        Some(TimeSeriesRowSelector::LastRow { after_merge: false })
1541                            if stream_ctx.input.sequence_range.is_none() =>
1542                        {
1543                            stream_ctx.input.series_row_selector
1544                        }
1545                        _ => None,
1546                    },
1547                    fetch_metrics.as_deref(),
1548                )
1549                .await?
1550            else {
1551                continue;
1552            };
1553            let build_cost = build_reader_start.elapsed();
1554            part_metrics.inc_build_reader_cost(build_cost);
1555
1556            let may_compat = range.compat_batch();
1557            let file_sequence_trusted = range
1558                .file_handle()
1559                .is_effective_target_sequence_trusted(stream_ctx.input.region_metadata().region_id);
1560
1561            let mapper = range.compaction_projection_mapper();
1562            while let Some(record_batch) = reader.next_batch().await? {
1563                let record_batch = if let Some(mapper) = mapper {
1564                    let batch = mapper.project(record_batch)?;
1565                    batch
1566                } else {
1567                    record_batch
1568                };
1569
1570                let Some(record_batch) = filter_flat_batch_by_sequence(
1571                    record_batch,
1572                    stream_ctx.input.sequence_range,
1573                    file_sequence_trusted,
1574                )? else {
1575                    continue;
1576                };
1577
1578                if let Some(flat_compat) = may_compat {
1579                    let batch = flat_compat.compat(record_batch)?;
1580                    yield batch;
1581                } else {
1582                    yield record_batch;
1583                }
1584            }
1585
1586            let prune_metrics = reader.metrics();
1587
1588            // Update per-file metrics if tracking is enabled
1589            if let Some(file_metrics_map) = per_file_metrics.as_mut() {
1590                let file_id = range.file_handle().file_id();
1591                let file_metrics = file_metrics_map
1592                    .entry(file_id)
1593                    .or_insert_with(FileScanMetrics::default);
1594
1595                file_metrics.num_ranges += 1;
1596                file_metrics.num_rows += prune_metrics.num_rows;
1597                file_metrics.build_reader_cost += build_cost;
1598                file_metrics.scan_cost += prune_metrics.scan_cost;
1599            }
1600
1601            reader_metrics.merge_from(&prune_metrics);
1602        }
1603
1604        // Reports metrics.
1605        reader_metrics.observe_rows(read_type);
1606        reader_metrics.filter_metrics.observe();
1607        part_metrics.merge_reader_metrics(reader_metrics, per_file_metrics.as_ref());
1608    }
1609}
1610
1611/// Build the stream of scanning the extension range in flat format denoted by the [`RowGroupIndex`].
1612#[cfg(feature = "enterprise")]
1613pub(crate) async fn scan_flat_extension_range(
1614    context: Arc<StreamContext>,
1615    index: RowGroupIndex,
1616    partition_metrics: PartitionMetrics,
1617    options: crate::extension::ExtensionRangeReadOptions,
1618) -> Result<BoxedRecordBatchStream> {
1619    use snafu::ResultExt;
1620
1621    let range = context.input.extension_range(index.index);
1622    let reader = range.flat_reader(context.as_ref(), options);
1623    let stream = reader
1624        .read(context, partition_metrics, index)
1625        .await
1626        .context(crate::error::ScanExternalRangeSnafu)?;
1627    Ok(stream)
1628}
1629
1630pub(crate) async fn maybe_scan_flat_other_ranges(
1631    context: &Arc<StreamContext>,
1632    index: RowGroupIndex,
1633    metrics: &PartitionMetrics,
1634    pre_filter_mode: PreFilterMode,
1635) -> Result<BoxedRecordBatchStream> {
1636    #[cfg(feature = "enterprise")]
1637    {
1638        let options = crate::extension::ExtensionRangeReadOptions { pre_filter_mode };
1639        scan_flat_extension_range(context.clone(), index, metrics.clone(), options).await
1640    }
1641
1642    #[cfg(not(feature = "enterprise"))]
1643    {
1644        let _ = context;
1645        let _ = index;
1646        let _ = metrics;
1647        let _ = pre_filter_mode;
1648
1649        crate::error::UnexpectedSnafu {
1650            reason: "no other ranges scannable in flat format",
1651        }
1652        .fail()
1653    }
1654}
1655
1656/// A stream wrapper that splits record batches from an inner stream.
1657pub(crate) struct SplitRecordBatchStream<S> {
1658    /// The inner stream that yields record batches.
1659    inner: S,
1660    /// Buffer for split batches.
1661    batches: VecDeque<RecordBatch>,
1662}
1663
1664impl<S> SplitRecordBatchStream<S> {
1665    /// Creates a new splitting stream wrapper.
1666    pub(crate) fn new(inner: S) -> Self {
1667        Self {
1668            inner,
1669            batches: VecDeque::new(),
1670        }
1671    }
1672}
1673
1674impl<S> Stream for SplitRecordBatchStream<S>
1675where
1676    S: Stream<Item = Result<RecordBatch>> + Unpin,
1677{
1678    type Item = Result<RecordBatch>;
1679
1680    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
1681        loop {
1682            // First, check if we have buffered split batches
1683            if let Some(batch) = self.batches.pop_front() {
1684                return Poll::Ready(Some(Ok(batch)));
1685            }
1686
1687            // Poll the inner stream for the next batch
1688            let record_batch = match futures::ready!(Pin::new(&mut self.inner).poll_next(cx)) {
1689                Some(Ok(batch)) => batch,
1690                Some(Err(e)) => return Poll::Ready(Some(Err(e))),
1691                None => return Poll::Ready(None),
1692            };
1693
1694            // Split the batch and buffer the results
1695            split_record_batch(record_batch, &mut self.batches);
1696            // Continue the loop to return the first split batch
1697        }
1698    }
1699}
1700
1701/// Splits the batch so each sub-batch has strictly increasing timestamps.
1702///
1703/// # Panics
1704/// Panics if the timestamp array is invalid.
1705pub(crate) fn split_record_batch(record_batch: RecordBatch, batches: &mut VecDeque<RecordBatch>) {
1706    let batch_rows = record_batch.num_rows();
1707    if batch_rows == 0 {
1708        return;
1709    }
1710    if batch_rows < 2 {
1711        batches.push_back(record_batch);
1712        return;
1713    }
1714
1715    let time_index_pos = time_index_column_index(record_batch.num_columns());
1716    let timestamps = record_batch.column(time_index_pos);
1717    let (ts_values, _unit) = timestamp_array_to_primitive(timestamps).unwrap();
1718    let mut offsets = Vec::with_capacity(16);
1719    offsets.push(0);
1720    let values = ts_values.values();
1721    for (i, &value) in values.iter().take(batch_rows - 1).enumerate() {
1722        if value >= values[i + 1] {
1723            offsets.push(i + 1);
1724        }
1725    }
1726    offsets.push(values.len());
1727
1728    // Splits the batch by offsets.
1729    for (i, &start) in offsets[..offsets.len() - 1].iter().enumerate() {
1730        let end = offsets[i + 1];
1731        let rows_in_batch = end - start;
1732        batches.push_back(record_batch.slice(start, rows_in_batch));
1733    }
1734}
1735
1736#[cfg(test)]
1737mod tests {
1738    use std::sync::Arc;
1739    use std::time::Instant;
1740
1741    use common_time::Timestamp;
1742    use smallvec::{SmallVec, smallvec};
1743    use store_api::storage::RegionId;
1744
1745    use super::*;
1746    use crate::cache::CacheStrategy;
1747    use crate::memtable::{
1748        BoxedBatchIterator, BoxedRecordBatchIterator, IterBuilder, MemtableRange,
1749        MemtableRangeContext, MemtableStats,
1750    };
1751    use crate::read::flat_projection::FlatProjectionMapper;
1752    use crate::read::range::{MemRangeBuilder, SourceIndex};
1753    use crate::read::scan_region::ScanInput;
1754    use crate::sst::file::{FileHandle, FileMeta};
1755    use crate::sst::file_purger::NoopFilePurger;
1756    use crate::test_util::memtable_util::metadata_for_test;
1757    use crate::test_util::scheduler_util::SchedulerEnv;
1758
1759    struct EmptyIterBuilder;
1760
1761    impl IterBuilder for EmptyIterBuilder {
1762        fn build(&self, _metrics: Option<MemScanMetrics>) -> Result<BoxedBatchIterator> {
1763            Ok(Box::new(std::iter::empty()))
1764        }
1765
1766        fn is_record_batch(&self) -> bool {
1767            true
1768        }
1769
1770        fn build_record_batch(
1771            &self,
1772            _time_range: Option<(Timestamp, Timestamp)>,
1773            _metrics: Option<MemScanMetrics>,
1774        ) -> Result<BoxedRecordBatchIterator> {
1775            Ok(Box::new(std::iter::empty()))
1776        }
1777    }
1778
1779    async fn new_test_stream_ctx(
1780        files: Vec<FileHandle>,
1781        memtables: Vec<MemRangeBuilder>,
1782    ) -> Arc<StreamContext> {
1783        let env = SchedulerEnv::new().await;
1784        let metadata = metadata_for_test();
1785        let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
1786        let input = ScanInput::builder(env.access_layer.clone(), mapper)
1787            .with_cache(CacheStrategy::Disabled)
1788            .with_memtables(memtables)
1789            .with_files(files)
1790            .build();
1791
1792        Arc::new(StreamContext {
1793            input,
1794            ranges: Vec::new(),
1795            query_start: Instant::now(),
1796        })
1797    }
1798
1799    fn new_test_file(num_rows: u64, num_series: u64) -> FileHandle {
1800        let meta = FileMeta {
1801            region_id: RegionId::new(123, 456),
1802            file_id: Default::default(),
1803            level: 1,
1804            time_range: (
1805                Timestamp::new_millisecond(0),
1806                Timestamp::new_millisecond(1000),
1807            ),
1808            num_rows,
1809            num_series,
1810            ..Default::default()
1811        };
1812        FileHandle::new(meta, Arc::new(NoopFilePurger))
1813    }
1814
1815    fn new_test_memtable(num_rows: usize, series_count: usize) -> MemRangeBuilder {
1816        let context = Arc::new(MemtableRangeContext::new(
1817            0,
1818            Box::new(EmptyIterBuilder),
1819            Default::default(),
1820        ));
1821        let stats = MemtableStats {
1822            time_range: Some((
1823                Timestamp::new_millisecond(0),
1824                Timestamp::new_millisecond(1000),
1825            )),
1826            num_rows,
1827            num_ranges: 1,
1828            series_count,
1829            ..Default::default()
1830        };
1831        let range = MemtableRange::new(context, stats.clone());
1832        MemRangeBuilder::new(range, stats)
1833    }
1834
1835    fn new_test_range_meta(row_group_indices: SmallVec<[RowGroupIndex; 2]>) -> RangeMeta {
1836        let indices = row_group_indices
1837            .iter()
1838            .map(|row_group_index| SourceIndex {
1839                index: row_group_index.index,
1840                num_row_groups: 1,
1841            })
1842            .collect();
1843
1844        RangeMeta {
1845            time_range: (
1846                Timestamp::new_millisecond(0),
1847                Timestamp::new_millisecond(1000),
1848            ),
1849            indices,
1850            row_group_indices,
1851            num_rows: 0,
1852        }
1853    }
1854
1855    #[tokio::test]
1856    async fn test_should_split_flat_batches_for_merge_uses_splittable_file_rows_per_series() {
1857        let num_rows = SPLIT_ROW_THRESHOLD * 2;
1858        let num_series = (num_rows / 100).max(1);
1859        let stream_ctx =
1860            new_test_stream_ctx(vec![new_test_file(num_rows, num_series)], vec![]).await;
1861        let range_meta = new_test_range_meta(smallvec![RowGroupIndex {
1862            index: 0,
1863            row_group_index: 0,
1864        }]);
1865
1866        assert_eq!(
1867            Some((num_rows / num_series) as usize),
1868            should_split_flat_batches_for_merge(&stream_ctx, &range_meta)
1869        );
1870    }
1871
1872    #[tokio::test]
1873    async fn test_should_split_flat_batches_for_merge_skips_small_or_unknown_series_files() {
1874        let stream_ctx = new_test_stream_ctx(
1875            vec![
1876                new_test_file(SPLIT_ROW_THRESHOLD.saturating_sub(1), 1),
1877                new_test_file(SPLIT_ROW_THRESHOLD * 2, 0),
1878            ],
1879            vec![],
1880        )
1881        .await;
1882        let range_meta = new_test_range_meta(smallvec![
1883            RowGroupIndex {
1884                index: 0,
1885                row_group_index: 0,
1886            },
1887            RowGroupIndex {
1888                index: 1,
1889                row_group_index: 0,
1890            }
1891        ]);
1892
1893        assert_eq!(
1894            None,
1895            should_split_flat_batches_for_merge(&stream_ctx, &range_meta)
1896        );
1897    }
1898
1899    #[tokio::test]
1900    async fn test_should_split_flat_batches_for_merge_returns_none_for_unsplittable_file() {
1901        let num_series =
1902            (SPLIT_ROW_THRESHOLD / (BATCH_SIZE_THRESHOLD - 1)).max(NUM_SERIES_THRESHOLD) + 1;
1903        let stream_ctx =
1904            new_test_stream_ctx(vec![new_test_file(SPLIT_ROW_THRESHOLD, num_series)], vec![]).await;
1905        let range_meta = new_test_range_meta(smallvec![RowGroupIndex {
1906            index: 0,
1907            row_group_index: 0,
1908        }]);
1909
1910        assert_eq!(
1911            None,
1912            should_split_flat_batches_for_merge(&stream_ctx, &range_meta)
1913        );
1914    }
1915
1916    #[tokio::test]
1917    async fn test_should_split_flat_batches_for_merge_falls_back_to_memtables() {
1918        let stream_ctx = new_test_stream_ctx(vec![], vec![new_test_memtable(5_000, 100)]).await;
1919        let range_meta = new_test_range_meta(smallvec![RowGroupIndex {
1920            index: 0,
1921            row_group_index: 0,
1922        }]);
1923
1924        assert_eq!(
1925            Some(50),
1926            should_split_flat_batches_for_merge(&stream_ctx, &range_meta)
1927        );
1928    }
1929
1930    #[tokio::test]
1931    async fn test_should_split_flat_batches_for_merge_clamps_estimate() {
1932        let stream_ctx =
1933            new_test_stream_ctx(vec![new_test_file(SPLIT_ROW_THRESHOLD * 2, 1)], vec![]).await;
1934        let range_meta = new_test_range_meta(smallvec![RowGroupIndex {
1935            index: 0,
1936            row_group_index: 0,
1937        }]);
1938
1939        assert_eq!(
1940            Some(DEFAULT_READ_BATCH_SIZE),
1941            should_split_flat_batches_for_merge(&stream_ctx, &range_meta)
1942        );
1943    }
1944
1945    #[test]
1946    fn test_compute_parallel_channel_size_clamps_to_max_for_small_batches() {
1947        assert_eq!(64, compute_parallel_channel_size(0));
1948        assert_eq!(64, compute_parallel_channel_size(1));
1949    }
1950
1951    #[test]
1952    fn test_compute_parallel_channel_size_returns_expected_mid_range_size() {
1953        assert_eq!(
1954            4,
1955            compute_parallel_channel_size(DEFAULT_READ_BATCH_SIZE / 2)
1956        );
1957    }
1958
1959    #[test]
1960    fn test_compute_parallel_channel_size_clamps_to_min_for_large_batches() {
1961        assert_eq!(2, compute_parallel_channel_size(DEFAULT_READ_BATCH_SIZE));
1962        assert_eq!(
1963            2,
1964            compute_parallel_channel_size(DEFAULT_READ_BATCH_SIZE * 2)
1965        );
1966    }
1967
1968    #[test]
1969    fn test_compute_average_batch_size_uses_arithmetic_mean() {
1970        assert_eq!(24, compute_average_batch_size([16, 24, 32]));
1971    }
1972
1973    #[test]
1974    fn test_compute_average_batch_size_clamps_values() {
1975        assert_eq!(
1976            DEFAULT_READ_BATCH_SIZE,
1977            compute_average_batch_size([DEFAULT_READ_BATCH_SIZE, DEFAULT_READ_BATCH_SIZE * 2])
1978        );
1979        assert_eq!(1, compute_average_batch_size([0, 1]));
1980    }
1981
1982    #[test]
1983    fn test_compute_average_batch_size_falls_back_when_empty() {
1984        assert_eq!(
1985            DEFAULT_READ_BATCH_SIZE,
1986            compute_average_batch_size(std::iter::empty())
1987        );
1988    }
1989
1990    /// Builds a flat-format record batch whose time index column holds `timestamps`.
1991    fn flat_ts_batch(timestamps: &[i64]) -> RecordBatch {
1992        use datatypes::arrow::array::{TimestampMillisecondArray, UInt8Array, UInt64Array};
1993        use datatypes::arrow::datatypes::{DataType, Field, Schema, TimeUnit};
1994
1995        let num_rows = timestamps.len();
1996        let schema = Arc::new(Schema::new(vec![
1997            Field::new(
1998                "ts",
1999                DataType::Timestamp(TimeUnit::Millisecond, None),
2000                false,
2001            ),
2002            Field::new("pk", DataType::UInt64, false),
2003            Field::new("seq", DataType::UInt64, false),
2004            Field::new("op", DataType::UInt8, false),
2005        ]));
2006        RecordBatch::try_new(
2007            schema,
2008            vec![
2009                Arc::new(TimestampMillisecondArray::from(timestamps.to_vec())),
2010                Arc::new(UInt64Array::from(vec![0u64; num_rows])),
2011                Arc::new(UInt64Array::from(vec![0u64; num_rows])),
2012                Arc::new(UInt8Array::from(vec![0u8; num_rows])),
2013            ],
2014        )
2015        .unwrap()
2016    }
2017
2018    /// Splits `timestamps` and returns the time index values of each sub-batch.
2019    fn split_ts(timestamps: &[i64]) -> Vec<Vec<i64>> {
2020        let mut batches = VecDeque::new();
2021        split_record_batch(flat_ts_batch(timestamps), &mut batches);
2022        batches
2023            .iter()
2024            .map(|batch| {
2025                let pos = time_index_column_index(batch.num_columns());
2026                let (values, _) = timestamp_array_to_primitive(batch.column(pos)).unwrap();
2027                values.values().to_vec()
2028            })
2029            .collect()
2030    }
2031
2032    #[test]
2033    fn test_split_record_batch_on_equal_timestamps() {
2034        // Splits on both decreasing and equal timestamps.
2035        assert_eq!(
2036            split_ts(&[1, 2, 2, 3, 1]),
2037            vec![vec![1, 2], vec![2, 3], vec![1]]
2038        );
2039        // A run of equal timestamps yields single-row sub-batches.
2040        assert_eq!(split_ts(&[5, 5, 5]), vec![vec![5], vec![5], vec![5]]);
2041        // Equal-ts run at the leading edge of the batch.
2042        assert_eq!(split_ts(&[5, 5, 1, 2]), vec![vec![5], vec![5], vec![1, 2]]);
2043        // Equal-ts run at the trailing edge of the batch.
2044        assert_eq!(split_ts(&[1, 2, 5, 5]), vec![vec![1, 2, 5], vec![5]]);
2045    }
2046
2047    #[test]
2048    fn test_split_record_batch_on_decreasing_timestamps() {
2049        assert_eq!(split_ts(&[1, 2, 3]), vec![vec![1, 2, 3]]);
2050        assert_eq!(split_ts(&[1, 3, 2, 4]), vec![vec![1, 3], vec![2, 4]]);
2051    }
2052
2053    #[test]
2054    fn test_split_record_batch_empty_and_single_row() {
2055        let mut batches = VecDeque::new();
2056        split_record_batch(flat_ts_batch(&[]), &mut batches);
2057        assert!(batches.is_empty());
2058
2059        assert_eq!(split_ts(&[42]), vec![vec![42]]);
2060    }
2061}
2062
2063#[cfg(test)]
2064mod sequence_filter_tests {
2065    use std::sync::Arc;
2066
2067    use datatypes::arrow::array::{Int64Array, StringArray, UInt8Array, UInt64Array};
2068    use datatypes::arrow::datatypes::{DataType, Field, Schema};
2069    use datatypes::arrow::record_batch::RecordBatch;
2070    use store_api::storage::SequenceRange;
2071
2072    use super::filter_flat_batch_by_sequence;
2073
2074    /// Builds a flat-format record batch: `(tag, field, ts, __primary_key, __sequence, __op_type)`.
2075    fn batch(sequences: &[u64]) -> RecordBatch {
2076        let schema = Arc::new(Schema::new(vec![
2077            Field::new("tag_0", DataType::Utf8, false),
2078            Field::new("field_0", DataType::Int64, false),
2079            Field::new(
2080                "ts",
2081                DataType::Timestamp(datatypes::arrow::datatypes::TimeUnit::Millisecond, None),
2082                false,
2083            ),
2084            Field::new("__primary_key", DataType::UInt8, false),
2085            Field::new("__sequence", DataType::UInt64, false),
2086            Field::new("__op_type", DataType::UInt8, false),
2087        ]));
2088        let tags = StringArray::from_iter_values((0..sequences.len()).map(|i| i.to_string()));
2089        let fields = Int64Array::from_iter_values(0..sequences.len() as i64);
2090        let ts = datatypes::arrow::array::TimestampMillisecondArray::from_iter_values(
2091            (0..sequences.len()).map(|i| i as i64 * 1000),
2092        );
2093        let pk = UInt8Array::from(vec![0u8; sequences.len()]);
2094        let seq = UInt64Array::from_iter_values(sequences.iter().copied());
2095        let op = UInt8Array::from(vec![0u8; sequences.len()]);
2096        RecordBatch::try_new(
2097            schema,
2098            vec![
2099                Arc::new(tags),
2100                Arc::new(fields),
2101                Arc::new(ts),
2102                Arc::new(pk),
2103                Arc::new(seq),
2104                Arc::new(op),
2105            ],
2106        )
2107        .unwrap()
2108    }
2109
2110    fn remaining_tags(batch: &RecordBatch) -> Vec<String> {
2111        batch
2112            .column(0)
2113            .as_any()
2114            .downcast_ref::<StringArray>()
2115            .unwrap()
2116            .iter()
2117            .map(|v| v.unwrap().to_string())
2118            .collect()
2119    }
2120
2121    #[test]
2122    fn test_filter_flat_batch_by_sequence_no_range_or_legacy_file() {
2123        let b = batch(&[1, 2, 3, 4]);
2124
2125        let out = filter_flat_batch_by_sequence(b.clone(), None, true).unwrap();
2126        assert_eq!(remaining_tags(&out.unwrap()), vec!["0", "1", "2", "3"]);
2127
2128        let out = filter_flat_batch_by_sequence(
2129            b.clone(),
2130            Some(SequenceRange::GtLtEq { min: 2, max: 3 }),
2131            false,
2132        )
2133        .unwrap();
2134        assert_eq!(remaining_tags(&out.unwrap()), vec!["0", "1", "2", "3"]);
2135    }
2136
2137    #[test]
2138    fn test_filter_flat_batch_by_sequence_exact_range() {
2139        let b = batch(&[1, 2, 3, 4]);
2140        let out = filter_flat_batch_by_sequence(
2141            b.clone(),
2142            Some(SequenceRange::GtLtEq { min: 2, max: 3 }),
2143            true,
2144        )
2145        .unwrap();
2146        assert_eq!(remaining_tags(&out.unwrap()), vec!["2"]);
2147
2148        let out = filter_flat_batch_by_sequence(
2149            b.clone(),
2150            Some(SequenceRange::GtLtEq { min: 10, max: 20 }),
2151            true,
2152        )
2153        .unwrap();
2154        assert!(out.is_none());
2155
2156        let empty = batch(&[]);
2157        let out = filter_flat_batch_by_sequence(
2158            empty,
2159            Some(SequenceRange::GtLtEq { min: 0, max: 10 }),
2160            true,
2161        )
2162        .unwrap();
2163        assert_eq!(out.unwrap().num_rows(), 0);
2164
2165        // Foreign batches are already virtualized by the reader. Their source
2166        // marker is irrelevant, and filtering uses the effective batch values.
2167        let out = filter_flat_batch_by_sequence(
2168            batch(&[1, 5, 9]),
2169            Some(SequenceRange::GtLtEq { min: 2, max: 8 }),
2170            true,
2171        )
2172        .unwrap();
2173        assert_eq!(remaining_tags(&out.unwrap()), vec!["1"]);
2174    }
2175}