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