1use 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#[derive(Default, Clone)]
59pub struct FileScanMetrics {
60 pub num_ranges: usize,
62 pub num_rows: usize,
64 pub build_part_cost: Duration,
66 pub build_reader_cost: Duration,
68 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 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#[derive(Default)]
110pub(crate) struct ScanMetricsSet {
111 prepare_scan_cost: Duration,
113 build_reader_cost: Duration,
115 scan_cost: Duration,
117 yield_cost: Duration,
119 convert_cost: Option<Time>,
121 total_cost: Duration,
123 num_rows: usize,
125 num_batches: usize,
127 num_mem_ranges: usize,
129 num_file_ranges: usize,
131
132 mem_scan_cost: Duration,
135 mem_rows: usize,
137 mem_batches: usize,
139 mem_series: usize,
141 mem_prefilter_cost: Duration,
143 mem_prefilter_rows_filtered: usize,
145
146 build_parts_cost: Duration,
149 sst_scan_cost: Duration,
151 rg_total: usize,
153 rg_fulltext_filtered: usize,
155 rg_inverted_filtered: usize,
157 rg_minmax_filtered: usize,
159 rg_bloom_filtered: usize,
161 rows_before_filter: usize,
163 rows_fulltext_filtered: usize,
165 rows_inverted_filtered: usize,
167 rows_bloom_filtered: usize,
169 rows_precise_filtered: usize,
171 fulltext_index_cache_hit: usize,
173 fulltext_index_cache_miss: usize,
175 inverted_index_cache_hit: usize,
177 inverted_index_cache_miss: usize,
179 bloom_filter_cache_hit: usize,
181 bloom_filter_cache_miss: usize,
183 minmax_cache_hit: usize,
185 minmax_cache_miss: usize,
187 pruner_cache_hit: usize,
189 pruner_cache_miss: usize,
191 pruner_prune_cost: Duration,
193 files_time_range_pruned: usize,
195 num_sst_record_batches: usize,
197 num_sst_batches: usize,
199 num_sst_rows: usize,
201
202 first_poll: Duration,
204
205 num_series_send_timeout: usize,
207 num_series_send_full: usize,
209 num_distributor_rows: usize,
211 num_distributor_batches: usize,
213 distributor_scan_cost: Duration,
215 distributor_yield_cost: Duration,
217 distributor_divider_cost: Duration,
219
220 merge_metrics: MergeMetrics,
222 dedup_metrics: DedupMetrics,
224
225 stream_eof: bool,
227
228 inverted_index_apply_metrics: Option<InvertedIndexApplyMetrics>,
231 bloom_filter_apply_metrics: Option<BloomFilterIndexApplyMetrics>,
233 fulltext_index_apply_metrics: Option<FulltextIndexApplyMetrics>,
235 fetch_metrics: Option<ParquetFetchMetrics>,
237 metadata_cache_metrics: Option<MetadataCacheMetrics>,
239 per_file_metrics: Option<HashMap<RegionFileId, FileScanMetrics>>,
241
242 build_ranges_mem_size: isize,
244 build_ranges_peak_mem_size: isize,
246 num_range_builders: isize,
248 num_peak_range_builders: isize,
250 range_cache_size: usize,
252 range_cache_hit: usize,
254 range_cache_miss: usize,
256}
257
258struct 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 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!(
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 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 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 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 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 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 if !merge_metrics.scan_cost.is_zero() {
545 write!(f, ", \"merge_metrics\":{:?}", merge_metrics)?;
546 }
547
548 if !dedup_metrics.dedup_cost.is_zero() {
550 write!(f, ", \"dedup_metrics\":{:?}", dedup_metrics)?;
551 }
552
553 if let Some(file_metrics) = per_file_metrics
555 && !file_metrics.is_empty()
556 {
557 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 total_cost.is_zero() && metrics.num_ranges == 0 {
566 continue;
567 }
568
569 if heap.len() < 10 {
570 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 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 fn with_prepare_scan_cost(mut self, cost: Duration) -> Self {
626 self.prepare_scan_cost += cost;
627 self
628 }
629
630 fn with_convert_cost(mut self, time: Time) -> Self {
632 self.convert_cost = Some(time);
633 self
634 }
635
636 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 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 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 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 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 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 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 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 partition: usize,
863 scanner_type: &'static str,
865 query_start: Instant,
867 explain_verbose: bool,
869 metrics: Mutex<ScanMetricsSet>,
871 in_progress_scan: IntGauge,
872
873 build_parts_cost: Time,
876 build_reader_cost: Time,
878 scan_cost: Time,
880 yield_cost: Time,
882 convert_cost: Time,
884 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 scan_metrics.merge_metrics.merge(metrics);
905
906 *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 scan_metrics.dedup_metrics.merge(metrics);
916
917 *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#[derive(Default)]
951pub(crate) struct PartitionMetricsList(Mutex<Vec<Option<PartitionMetrics>>>);
952
953impl PartitionMetricsList {
954 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 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#[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 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 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 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 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 if let Some(file_metrics) = per_file_metrics {
1084 metrics_set.merge_per_file_metrics(file_metrics);
1085 }
1086 }
1087
1088 pub(crate) fn on_finish(&self) {
1090 self.0.on_finish(true);
1091 }
1092
1093 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 pub(crate) fn explain_verbose(&self) -> bool {
1101 self.0.explain_verbose
1102 }
1103
1104 pub(crate) fn merge_metrics_reporter(&self) -> Arc<dyn MergeMetricsReport> {
1106 self.0.clone()
1107 }
1108
1109 pub(crate) fn dedup_metrics_reporter(&self) -> Arc<dyn DedupMetricsReport> {
1111 self.0.clone()
1112 }
1113
1114 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 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 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#[derive(Default)]
1146pub(crate) struct SeriesDistributorMetrics {
1147 pub(crate) num_series_send_timeout: usize,
1149 pub(crate) num_series_send_full: usize,
1151 pub(crate) num_rows: usize,
1153 pub(crate) num_batches: usize,
1155 pub(crate) scan_cost: Duration,
1157 pub(crate) yield_cost: Duration,
1159 pub(crate) divider_cost: Duration,
1161}
1162
1163#[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 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
1200const SPLIT_ROW_THRESHOLD: u64 = DEFAULT_ROW_GROUP_SIZE as u64;
1202const NUM_SERIES_THRESHOLD: u64 = 10240;
1204const BATCH_SIZE_THRESHOLD: u64 = 50;
1207
1208pub(crate) fn should_split_flat_batches_for_merge(
1211 stream_ctx: &Arc<StreamContext>,
1212 range_meta: &RangeMeta,
1213) -> Option<usize> {
1214 let mut num_files_to_split = 0;
1216 let mut num_mem_rows = 0;
1217 let mut num_mem_series = 0;
1218 let mut total_rows: u64 = 0;
1220 let mut total_series: u64 = 0;
1221 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 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 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 num_files_to_split += 1;
1239 continue;
1240 } else if file_meta.num_rows < SPLIT_ROW_THRESHOLD || file_meta.num_series == 0 {
1241 continue;
1243 }
1244 debug_assert!(file_meta.num_rows > 0);
1245 if !can_split_series(file_meta.num_rows, file_meta.num_series) {
1246 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 }
1263
1264 let should_split = if num_files_to_split > 0 {
1265 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 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 DEFAULT_READ_BATCH_SIZE / 4
1288 };
1289 Some(estimated_batch_size)
1290}
1291
1292pub(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
1299pub(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 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
1402pub(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#[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 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
1470pub(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#[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 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 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#[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
1656pub(crate) struct SplitRecordBatchStream<S> {
1658 inner: S,
1660 batches: VecDeque<RecordBatch>,
1662}
1663
1664impl<S> SplitRecordBatchStream<S> {
1665 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 if let Some(batch) = self.batches.pop_front() {
1684 return Poll::Ready(Some(Ok(batch)));
1685 }
1686
1687 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_record_batch(record_batch, &mut self.batches);
1696 }
1698 }
1699}
1700
1701pub(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 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 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 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 assert_eq!(
2036 split_ts(&[1, 2, 2, 3, 1]),
2037 vec![vec![1, 2], vec![2, 3], vec![1]]
2038 );
2039 assert_eq!(split_ts(&[5, 5, 5]), vec![vec![5], vec![5], vec![5]]);
2041 assert_eq!(split_ts(&[5, 5, 1, 2]), vec![vec![5], vec![5], vec![1, 2]]);
2043 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 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 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}