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};
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 rg_vector_filtered: usize,
163 rows_before_filter: usize,
165 rows_fulltext_filtered: usize,
167 rows_inverted_filtered: usize,
169 rows_bloom_filtered: usize,
171 rows_vector_filtered: usize,
173 rows_vector_selected: usize,
175 rows_precise_filtered: usize,
177 fulltext_index_cache_hit: usize,
179 fulltext_index_cache_miss: usize,
181 inverted_index_cache_hit: usize,
183 inverted_index_cache_miss: usize,
185 bloom_filter_cache_hit: usize,
187 bloom_filter_cache_miss: usize,
189 minmax_cache_hit: usize,
191 minmax_cache_miss: usize,
193 pruner_cache_hit: usize,
195 pruner_cache_miss: usize,
197 pruner_prune_cost: Duration,
199 files_time_range_pruned: usize,
201 num_sst_record_batches: usize,
203 num_sst_batches: usize,
205 num_sst_rows: usize,
207
208 first_poll: Duration,
210
211 num_series_send_timeout: usize,
213 num_series_send_full: usize,
215 num_distributor_rows: usize,
217 num_distributor_batches: usize,
219 distributor_scan_cost: Duration,
221 distributor_yield_cost: Duration,
223 distributor_divider_cost: Duration,
225
226 merge_metrics: MergeMetrics,
228 dedup_metrics: DedupMetrics,
230
231 stream_eof: bool,
233
234 inverted_index_apply_metrics: Option<InvertedIndexApplyMetrics>,
237 bloom_filter_apply_metrics: Option<BloomFilterIndexApplyMetrics>,
239 fulltext_index_apply_metrics: Option<FulltextIndexApplyMetrics>,
241 fetch_metrics: Option<ParquetFetchMetrics>,
243 metadata_cache_metrics: Option<MetadataCacheMetrics>,
245 per_file_metrics: Option<HashMap<RegionFileId, FileScanMetrics>>,
247
248 build_ranges_mem_size: isize,
250 build_ranges_peak_mem_size: isize,
252 num_range_builders: isize,
254 num_peak_range_builders: isize,
256 range_cache_size: usize,
258 range_cache_hit: usize,
260 range_cache_miss: usize,
262}
263
264struct 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 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!(
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 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 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 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 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 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 if !merge_metrics.scan_cost.is_zero() {
563 write!(f, ", \"merge_metrics\":{:?}", merge_metrics)?;
564 }
565
566 if !dedup_metrics.dedup_cost.is_zero() {
568 write!(f, ", \"dedup_metrics\":{:?}", dedup_metrics)?;
569 }
570
571 if let Some(file_metrics) = per_file_metrics
573 && !file_metrics.is_empty()
574 {
575 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 total_cost.is_zero() && metrics.num_ranges == 0 {
584 continue;
585 }
586
587 if heap.len() < 10 {
588 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 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 fn with_prepare_scan_cost(mut self, cost: Duration) -> Self {
644 self.prepare_scan_cost += cost;
645 self
646 }
647
648 fn with_convert_cost(mut self, time: Time) -> Self {
650 self.convert_cost = Some(time);
651 self
652 }
653
654 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 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 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 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 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 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 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 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 partition: usize,
895 scanner_type: &'static str,
897 query_start: Instant,
899 explain_verbose: bool,
901 metrics: Mutex<ScanMetricsSet>,
903 in_progress_scan: IntGauge,
904
905 build_parts_cost: Time,
908 build_reader_cost: Time,
910 scan_cost: Time,
912 yield_cost: Time,
914 convert_cost: Time,
916 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 scan_metrics.merge_metrics.merge(metrics);
937
938 *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 scan_metrics.dedup_metrics.merge(metrics);
948
949 *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#[derive(Default)]
983pub(crate) struct PartitionMetricsList(Mutex<Vec<Option<PartitionMetrics>>>);
984
985impl PartitionMetricsList {
986 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 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#[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 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 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 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 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 if let Some(file_metrics) = per_file_metrics {
1116 metrics_set.merge_per_file_metrics(file_metrics);
1117 }
1118 }
1119
1120 pub(crate) fn on_finish(&self) {
1122 self.0.on_finish(true);
1123 }
1124
1125 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 pub(crate) fn explain_verbose(&self) -> bool {
1133 self.0.explain_verbose
1134 }
1135
1136 pub(crate) fn merge_metrics_reporter(&self) -> Arc<dyn MergeMetricsReport> {
1138 self.0.clone()
1139 }
1140
1141 pub(crate) fn dedup_metrics_reporter(&self) -> Arc<dyn DedupMetricsReport> {
1143 self.0.clone()
1144 }
1145
1146 #[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 #[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 #[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#[derive(Default)]
1181pub(crate) struct SeriesDistributorMetrics {
1182 pub(crate) num_series_send_timeout: usize,
1184 pub(crate) num_series_send_full: usize,
1186 pub(crate) num_rows: usize,
1188 pub(crate) num_batches: usize,
1190 pub(crate) scan_cost: Duration,
1192 pub(crate) yield_cost: Duration,
1194 pub(crate) divider_cost: Duration,
1196}
1197
1198#[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 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
1235const SPLIT_ROW_THRESHOLD: u64 = DEFAULT_ROW_GROUP_SIZE as u64;
1237const NUM_SERIES_THRESHOLD: u64 = 10240;
1239const BATCH_SIZE_THRESHOLD: u64 = 50;
1242
1243pub(crate) fn should_split_flat_batches_for_merge(
1246 stream_ctx: &Arc<StreamContext>,
1247 range_meta: &RangeMeta,
1248) -> Option<usize> {
1249 let mut num_files_to_split = 0;
1251 let mut num_mem_rows = 0;
1252 let mut num_mem_series = 0;
1253 let mut total_rows: u64 = 0;
1255 let mut total_series: u64 = 0;
1256 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 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 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 num_files_to_split += 1;
1274 continue;
1275 } else if file_meta.num_rows < SPLIT_ROW_THRESHOLD || file_meta.num_series == 0 {
1276 continue;
1278 }
1279 debug_assert!(file_meta.num_rows > 0);
1280 if !can_split_series(file_meta.num_rows, file_meta.num_series) {
1281 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 }
1298
1299 let should_split = if num_files_to_split > 0 {
1300 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 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 DEFAULT_READ_BATCH_SIZE / 4
1323 };
1324 Some(estimated_batch_size)
1325}
1326
1327pub(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
1334pub(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 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
1437pub(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#[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 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
1505fn 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#[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 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 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 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#[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
1696pub(crate) struct SplitRecordBatchStream<S> {
1698 inner: S,
1700 batches: VecDeque<RecordBatch>,
1702}
1703
1704impl<S> SplitRecordBatchStream<S> {
1705 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 if let Some(batch) = self.batches.pop_front() {
1724 return Poll::Ready(Some(Ok(batch)));
1725 }
1726
1727 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_record_batch(record_batch, &mut self.batches);
1736 }
1738 }
1739}
1740
1741pub(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 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 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 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 assert_eq!(
2076 split_ts(&[1, 2, 2, 3, 1]),
2077 vec![vec![1, 2], vec![2, 3], vec![1]]
2078 );
2079 assert_eq!(split_ts(&[5, 5, 5]), vec![vec![5], vec![5], vec![5]]);
2081 assert_eq!(split_ts(&[5, 5, 1, 2]), vec![vec![5], vec![5], vec![1, 2]]);
2083 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 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 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}