1use std::collections::{HashMap, HashSet};
18use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
19use std::sync::{Arc, Mutex};
20use std::time::Instant;
21
22use common_telemetry::debug;
23use smallvec::SmallVec;
24use snafu::ResultExt;
25use store_api::region_engine::PartitionRange;
26use store_api::storage::FileId;
27use tokio::sync::{mpsc, oneshot};
28use uuid::Uuid;
29
30use crate::error::{PruneFileSnafu, Result};
31use crate::metrics::PRUNER_ACTIVE_BUILDERS;
32use crate::read::range::{FileRangeBuilder, RowGroupIndex};
33use crate::read::scan_region::StreamContext;
34use crate::read::scan_util::{FileScanMetrics, PartitionMetrics, new_filter_metrics};
35use crate::sst::parquet::file_range::{FileRange, PreFilterMode};
36use crate::sst::parquet::reader::ReaderMetrics;
37
38const PREFETCH_COUNT: usize = 8;
40
41pub struct PartitionPruner {
43 pruner: Arc<Pruner>,
44 file_indices: Vec<usize>,
46 pre_filter_modes: Vec<PreFilterMode>,
48 current_position: AtomicUsize,
50}
51
52impl PartitionPruner {
53 pub fn new(pruner: Arc<Pruner>, partition_ranges: &[PartitionRange]) -> Self {
55 let num_files = pruner.inner.stream_ctx.input.num_files();
56 let mut file_indices = Vec::with_capacity(num_files);
57 let mut pre_filter_modes = vec![PreFilterMode::SkipFields; num_files];
58 let mut dedup_set = HashSet::with_capacity(pruner.inner.stream_ctx.input.num_files());
59
60 let num_memtables = pruner.inner.stream_ctx.input.num_memtables();
61 for part_range in partition_ranges {
62 let range_meta = &pruner.inner.stream_ctx.ranges[part_range.identifier];
63 let pre_filter_mode = pruner.inner.stream_ctx.range_pre_filter_mode(part_range);
64 for row_group_index in &range_meta.row_group_indices {
65 if pruner
66 .inner
67 .stream_ctx
68 .is_file_range_index(*row_group_index)
69 {
70 let file_index = row_group_index.index - num_memtables;
71 if dedup_set.contains(&file_index) {
72 continue;
73 } else {
74 file_indices.push(file_index);
75 pre_filter_modes[file_index] = pre_filter_mode;
76 dedup_set.insert(file_index);
77 }
78 }
79 }
80 }
81
82 Self {
83 pruner,
84 file_indices,
85 pre_filter_modes,
86 current_position: AtomicUsize::new(0),
87 }
88 }
89
90 pub async fn build_file_ranges(
95 &self,
96 index: RowGroupIndex,
97 partition_metrics: &PartitionMetrics,
98 reader_metrics: &mut ReaderMetrics,
99 ) -> Result<SmallVec<[FileRange; 2]>> {
100 let file_index = index.index - self.pruner.inner.stream_ctx.input.num_memtables();
101 let pre_filter_mode = self.pre_filter_mode(file_index);
102
103 let ranges = self
105 .pruner
106 .build_file_ranges(index, pre_filter_mode, partition_metrics, reader_metrics)
107 .await?;
108
109 if let Some(pos) = self.file_indices.iter().position(|&idx| idx == file_index) {
111 let prev_pos = self.current_position.fetch_max(pos, Ordering::Relaxed);
112 if pos > prev_pos || prev_pos == 0 {
113 self.prefetch_upcoming_files(pos, partition_metrics);
114 }
115 }
116
117 Ok(ranges)
118 }
119
120 pub fn try_skip_manifest_pruned_file_range(
130 &self,
131 index: RowGroupIndex,
132 part_metrics: &PartitionMetrics,
133 ) -> bool {
134 let Some(file_index) = self.file_index(index) else {
135 return false;
136 };
137 let mut reader_metrics = ReaderMetrics::default();
138 let pruned = self
139 .pruner
140 .inner
141 .try_mark_manifest_pruned(file_index, &mut reader_metrics);
142 if pruned {
143 self.pruner.skip_file_range(index, &mut reader_metrics);
144 part_metrics.merge_reader_metrics(&reader_metrics, None);
145 }
146 pruned
147 }
148
149 fn prefetch_upcoming_files(&self, current_pos: usize, partition_metrics: &PartitionMetrics) {
151 let start = current_pos + 1;
152 let end = (start + PREFETCH_COUNT).min(self.file_indices.len());
153
154 for i in start..end {
155 let file_index = self.file_indices[i];
156 let pre_filter_mode = self.pre_filter_mode(file_index);
157 self.pruner.get_file_builder_background(
158 file_index,
159 pre_filter_mode,
160 Some(partition_metrics.clone()),
161 );
162 }
163 }
164
165 fn pre_filter_mode(&self, file_index: usize) -> PreFilterMode {
166 self.pre_filter_modes
167 .get(file_index)
168 .copied()
169 .unwrap_or(PreFilterMode::SkipFields)
170 }
171
172 fn file_index(&self, index: RowGroupIndex) -> Option<usize> {
173 self.pruner
174 .inner
175 .stream_ctx
176 .is_file_range_index(index)
177 .then(|| index.index - self.pruner.inner.stream_ctx.input.num_memtables())
178 }
179}
180
181#[derive(Debug, Clone, Copy)]
183pub struct PrunerOptions {
184 pub retain_builders: bool,
186 pub enable_predicate_prefilter: bool,
194}
195
196impl Default for PrunerOptions {
197 fn default() -> Self {
198 Self {
199 retain_builders: false,
200 enable_predicate_prefilter: true,
201 }
202 }
203}
204
205pub struct Pruner {
207 worker_senders: Vec<mpsc::Sender<PruneRequest>>,
209 inner: Arc<PrunerInner>,
210}
211
212struct PrunerInner {
213 num_workers: usize,
215 file_entries: Vec<Mutex<FileBuilderEntry>>,
217 stream_ctx: Arc<StreamContext>,
219 manifest_pruned_files: Vec<AtomicBool>,
225 retain_builders: bool,
227 enable_predicate_prefilter: bool,
229}
230
231impl Drop for PrunerInner {
232 fn drop(&mut self) {
233 let active_builders = self
234 .file_entries
235 .iter_mut()
236 .map(|entry| usize::from(entry.get_mut().is_ok_and(|entry| entry.builder.is_some())))
237 .sum::<usize>();
238 PRUNER_ACTIVE_BUILDERS.sub(active_builders as i64);
239 }
240}
241
242impl PrunerInner {
243 fn try_mark_manifest_pruned(
249 &self,
250 file_index: usize,
251 reader_metrics: &mut ReaderMetrics,
252 ) -> bool {
253 if self.manifest_pruned_files[file_index].load(Ordering::Relaxed) {
254 return true;
255 }
256 let file = &self.stream_ctx.input.files[file_index];
257 if !self.stream_ctx.input.can_manifest_prune_file(file) {
258 return false;
259 }
260 if self.manifest_pruned_files[file_index]
261 .compare_exchange(false, true, Ordering::Relaxed, Ordering::Relaxed)
262 .is_ok()
263 {
264 reader_metrics.filter_metrics.files_time_range_pruned += 1;
265 }
266 true
267 }
268}
269
270struct FileBuilderEntry {
272 builder: Option<Arc<FileRangeBuilder>>,
278 remaining_ranges: usize,
282 waiters: Vec<oneshot::Sender<Result<Arc<FileRangeBuilder>>>>,
284}
285
286struct PruneRequest {
288 file_index: usize,
290 pre_filter_mode: PreFilterMode,
292 response_tx: Option<oneshot::Sender<Result<Arc<FileRangeBuilder>>>>,
294 partition_metrics: Option<PartitionMetrics>,
296}
297
298impl Pruner {
299 pub fn new(stream_ctx: Arc<StreamContext>, num_workers: usize) -> Self {
304 Self::new_with_options(stream_ctx, num_workers, PrunerOptions::default())
305 }
306
307 pub fn new_with_options(
309 stream_ctx: Arc<StreamContext>,
310 num_workers: usize,
311 options: PrunerOptions,
312 ) -> Self {
313 let PrunerOptions {
314 retain_builders,
315 enable_predicate_prefilter,
316 } = options;
317 let num_files = stream_ctx.input.num_files();
318 let file_entries: Vec<_> = (0..num_files)
319 .map(|_| {
320 Mutex::new(FileBuilderEntry {
321 builder: None,
322 remaining_ranges: 0,
323 waiters: Vec::new(),
324 })
325 })
326 .collect();
327 let manifest_pruned_files: Vec<AtomicBool> =
328 (0..num_files).map(|_| AtomicBool::new(false)).collect();
329 let mut worker_senders = Vec::with_capacity(num_workers);
331 let mut receivers = Vec::with_capacity(num_workers);
332 for _ in 0..num_workers {
333 let (tx, rx) = mpsc::channel::<PruneRequest>(64);
334 worker_senders.push(tx);
335 receivers.push(rx);
336 }
337
338 let inner = Arc::new(PrunerInner {
339 num_workers,
340 file_entries,
341 stream_ctx,
342 manifest_pruned_files,
343 retain_builders,
344 enable_predicate_prefilter,
345 });
346
347 for (worker_id, rx) in receivers.into_iter().enumerate() {
349 let inner_clone = inner.clone();
350 common_runtime::spawn_query(async move {
351 Self::worker_loop(worker_id, rx, inner_clone).await;
352 });
353 }
354
355 Self {
356 worker_senders,
357 inner,
358 }
359 }
360
361 pub fn add_partition_ranges(&self, partition_ranges: &[PartitionRange]) {
368 for pruned in &self.inner.manifest_pruned_files {
370 pruned.store(false, Ordering::Relaxed);
371 }
372
373 let num_memtables = self.inner.stream_ctx.input.num_memtables();
375 for part_range in partition_ranges {
376 let range_meta = &self.inner.stream_ctx.ranges[part_range.identifier];
377 for row_group_index in &range_meta.row_group_indices {
378 if self.inner.stream_ctx.is_file_range_index(*row_group_index) {
379 let file_index = row_group_index.index - num_memtables;
380 if file_index < self.inner.file_entries.len() {
381 let mut entry = self.inner.file_entries[file_index].lock().unwrap();
382 entry.remaining_ranges += 1;
383 }
384 }
385 }
386 }
387 }
388
389 pub async fn build_file_ranges(
395 &self,
396 index: RowGroupIndex,
397 pre_filter_mode: PreFilterMode,
398 partition_metrics: &PartitionMetrics,
399 reader_metrics: &mut ReaderMetrics,
400 ) -> Result<SmallVec<[FileRange; 2]>> {
401 let file_index = index.index - self.inner.stream_ctx.input.num_memtables();
402
403 let builder = self
405 .get_file_builder(
406 file_index,
407 pre_filter_mode,
408 partition_metrics,
409 reader_metrics,
410 )
411 .await?;
412
413 let mut ranges = SmallVec::new();
415 builder.build_ranges(index.row_group_index, &mut ranges);
416
417 self.decrement_and_maybe_clear(file_index, reader_metrics);
419
420 Ok(ranges)
421 }
422
423 pub fn skip_file_range(&self, index: RowGroupIndex, reader_metrics: &mut ReaderMetrics) {
429 if !self.inner.stream_ctx.is_file_range_index(index) {
430 return;
431 }
432 let file_index = index.index - self.inner.stream_ctx.input.num_memtables();
433 self.decrement_and_maybe_clear(file_index, reader_metrics);
434 }
435
436 async fn get_file_builder(
438 &self,
439 file_index: usize,
440 pre_filter_mode: PreFilterMode,
441 partition_metrics: &PartitionMetrics,
442 reader_metrics: &mut ReaderMetrics,
443 ) -> Result<Arc<FileRangeBuilder>> {
444 {
446 let entry = self.inner.file_entries[file_index].lock().unwrap();
447 if let Some(builder) = &entry.builder {
448 reader_metrics.filter_metrics.pruner_cache_hit += 1;
449 return Ok(builder.clone());
450 }
451 }
452
453 reader_metrics.filter_metrics.pruner_cache_miss += 1;
454 let prune_start = Instant::now();
455 let file = &self.inner.stream_ctx.input.files[file_index];
456 let file_id = file.file_id().file_id();
457 let worker_idx = self.get_worker_idx(file_id);
458
459 let (response_tx, response_rx) = oneshot::channel();
460 let request = PruneRequest {
461 file_index,
462 pre_filter_mode,
463 response_tx: Some(response_tx),
464 partition_metrics: Some(partition_metrics.clone()),
465 };
466
467 let result = if self.worker_senders[worker_idx].send(request).await.is_err() {
468 common_telemetry::warn!("Worker channel closed, falling back to direct pruning");
469 self.prune_file_directly(file_index, pre_filter_mode, reader_metrics)
471 .await
472 } else {
473 match response_rx.await {
475 Ok(result) => result,
476 Err(_) => {
477 common_telemetry::warn!(
478 "Response channel closed, falling back to direct pruning"
479 );
480 self.prune_file_directly(file_index, pre_filter_mode, reader_metrics)
482 .await
483 }
484 }
485 };
486 reader_metrics.filter_metrics.pruner_prune_cost += prune_start.elapsed();
487 result
488 }
489
490 pub fn get_file_builder_background(
492 &self,
493 file_index: usize,
494 pre_filter_mode: PreFilterMode,
495 partition_metrics: Option<PartitionMetrics>,
496 ) {
497 {
499 let entry = self.inner.file_entries[file_index].lock().unwrap();
500 if entry.builder.is_some() {
501 return;
502 }
503 }
504
505 let file = &self.inner.stream_ctx.input.files[file_index];
506 let file_id = file.file_id().file_id();
507 let worker_idx = self.get_worker_idx(file_id);
508
509 let request = PruneRequest {
510 file_index,
511 pre_filter_mode,
512 response_tx: None,
513 partition_metrics,
514 };
515
516 let _ = self.worker_senders[worker_idx].try_send(request);
518 }
519
520 pub fn predicate_prefilter_enabled(&self) -> bool {
523 self.inner.enable_predicate_prefilter
524 }
525
526 fn get_worker_idx(&self, file_id: FileId) -> usize {
527 let file_id_hash = Uuid::from(file_id).as_u128() as usize;
528 file_id_hash % self.inner.num_workers
529 }
530
531 async fn prune_file_directly(
534 &self,
535 file_index: usize,
536 pre_filter_mode: PreFilterMode,
537 reader_metrics: &mut ReaderMetrics,
538 ) -> Result<Arc<FileRangeBuilder>> {
539 if self
541 .inner
542 .try_mark_manifest_pruned(file_index, reader_metrics)
543 {
544 let arc_builder = Arc::new(FileRangeBuilder::default());
545 return Ok(arc_builder);
548 }
549
550 let file = &self.inner.stream_ctx.input.files[file_index];
551 let predicate = self.inner.stream_ctx.input.predicate_for_file(file);
552 let builder = self
553 .inner
554 .stream_ctx
555 .input
556 .prune_file_after_manifest_check(
557 file,
558 pre_filter_mode,
559 self.inner.enable_predicate_prefilter,
560 predicate,
561 reader_metrics,
562 )
563 .await?;
564
565 let arc_builder = Arc::new(builder);
566
567 {
570 let mut entry = self.inner.file_entries[file_index].lock().unwrap();
571 cache_builder_if_needed(&mut entry, &arc_builder, reader_metrics);
572 }
573
574 Ok(arc_builder)
575 }
576
577 fn decrement_and_maybe_clear(&self, file_index: usize, reader_metrics: &mut ReaderMetrics) {
579 let mut entry = self.inner.file_entries[file_index].lock().unwrap();
580 entry.remaining_ranges = entry.remaining_ranges.saturating_sub(1);
581
582 if !self.inner.retain_builders
583 && entry.remaining_ranges == 0
584 && let Some(builder) = entry.builder.take()
585 {
586 PRUNER_ACTIVE_BUILDERS.dec();
587 reader_metrics.metadata_mem_size -= builder.memory_size() as isize;
588 reader_metrics.num_range_builders -= 1;
589 }
590 }
591
592 async fn worker_loop(
594 worker_id: usize,
595 mut rx: mpsc::Receiver<PruneRequest>,
596 inner: Arc<PrunerInner>,
597 ) {
598 let mut worker_cache_hit = 0;
599 let mut worker_cache_miss = 0;
600 let mut pruned_files = Vec::new();
601
602 while let Some(request) = rx.recv().await {
603 let PruneRequest {
604 file_index,
605 pre_filter_mode,
606 response_tx,
607 partition_metrics,
608 } = request;
609
610 {
612 let entry = inner.file_entries[file_index].lock().unwrap();
613 if let Some(builder) = &entry.builder {
614 if let Some(response_tx) = response_tx {
616 let _ = response_tx.send(Ok(builder.clone()));
617 }
618 worker_cache_hit += 1;
619 continue;
620 }
621 }
622 worker_cache_miss += 1;
623
624 let file = &inner.stream_ctx.input.files[file_index];
625 pruned_files.push(file.file_id().file_id());
626 let explain_verbose = partition_metrics
627 .as_ref()
628 .map(|m| m.explain_verbose())
629 .unwrap_or(false);
630 let mut metrics = ReaderMetrics {
631 filter_metrics: new_filter_metrics(explain_verbose),
632 ..Default::default()
633 };
634
635 let result = if inner.try_mark_manifest_pruned(file_index, &mut metrics) {
637 Ok(FileRangeBuilder::default())
640 } else {
641 let predicate = inner.stream_ctx.input.predicate_for_file(file);
642 inner
643 .stream_ctx
644 .input
645 .prune_file_after_manifest_check(
646 file,
647 pre_filter_mode,
648 inner.enable_predicate_prefilter,
649 predicate,
650 &mut metrics,
651 )
652 .await
653 };
654
655 let mut entry = inner.file_entries[file_index].lock().unwrap();
657 match result {
658 Ok(builder) => {
659 let arc_builder = Arc::new(builder);
660 let is_background = response_tx.is_none();
661
662 let did_cache =
668 if inner.manifest_pruned_files[file_index].load(Ordering::Relaxed) {
669 false
670 } else {
671 cache_builder_if_needed(&mut entry, &arc_builder, &mut metrics)
672 };
673
674 for waiter in entry.waiters.drain(..) {
676 let _ = waiter.send(Ok(arc_builder.clone()));
677 }
678 if let Some(response_tx) = response_tx {
680 let _ = response_tx.send(Ok(arc_builder));
681 }
682
683 debug!(
684 "Pruner worker {} pruned file_index: {}, file: {:?}, metrics: {:?}",
685 worker_id,
686 file_index,
687 file.file_id(),
688 metrics
689 );
690
691 if (!is_background || did_cache)
696 && let Some(part_metrics) = &partition_metrics
697 {
698 let per_file_metrics = if part_metrics.explain_verbose() {
699 let file_id = file.file_id();
700 let mut map = HashMap::new();
701 map.insert(
702 file_id,
703 FileScanMetrics {
704 build_part_cost: metrics.build_cost,
705 ..Default::default()
706 },
707 );
708 Some(map)
709 } else {
710 None
711 };
712 part_metrics.merge_reader_metrics(&metrics, per_file_metrics.as_ref());
713 }
714 }
715 Err(e) => {
716 let arc_error = Arc::new(e);
717 for waiter in entry.waiters.drain(..) {
718 let _ = waiter.send(Err(arc_error.clone()).context(PruneFileSnafu));
719 }
720 if let Some(response_tx) = response_tx {
721 let _ = response_tx.send(Err(arc_error).context(PruneFileSnafu));
722 }
723 }
724 }
725 }
726
727 common_telemetry::debug!(
728 "Pruner worker {} finished, cache_hit: {}, cache_miss: {}, files: {:?}",
729 worker_id,
730 worker_cache_hit,
731 worker_cache_miss,
732 pruned_files,
733 );
734 }
735}
736
737#[cfg(test)]
738impl Pruner {
739 fn test_remaining_ranges(&self, file_index: usize) -> usize {
741 self.inner.file_entries[file_index]
742 .lock()
743 .unwrap()
744 .remaining_ranges
745 }
746
747 fn test_has_builder(&self, file_index: usize) -> bool {
749 self.inner.file_entries[file_index]
750 .lock()
751 .unwrap()
752 .builder
753 .is_some()
754 }
755
756 fn test_is_manifest_pruned(&self, file_index: usize) -> bool {
758 self.inner.manifest_pruned_files[file_index].load(Ordering::Relaxed)
759 }
760
761 #[allow(dead_code)]
763 fn test_clear_builder(&self, file_index: usize) {
764 let mut entry = self.inner.file_entries[file_index].lock().unwrap();
765 if entry.builder.take().is_some() {
766 PRUNER_ACTIVE_BUILDERS.dec();
767 }
768 }
769}
770
771fn should_cache_builder(entry: &FileBuilderEntry) -> bool {
773 entry.builder.is_none() && entry.remaining_ranges > 0
774}
775
776fn cache_builder_if_needed(
779 entry: &mut FileBuilderEntry,
780 builder: &Arc<FileRangeBuilder>,
781 reader_metrics: &mut ReaderMetrics,
782) -> bool {
783 if should_cache_builder(entry) {
784 reader_metrics.metadata_mem_size += builder.memory_size() as isize;
785 reader_metrics.num_range_builders += 1;
786 entry.builder = Some(builder.clone());
787 PRUNER_ACTIVE_BUILDERS.inc();
788 true
789 } else {
790 false
791 }
792}
793
794#[cfg(test)]
795mod tests {
796 use common_time::Timestamp;
797 use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet;
798 use datafusion_common::ScalarValue;
799 use datafusion_expr::{Expr, col, lit};
800 use store_api::region_engine::PartitionRange;
801 use store_api::storage::{FileId, RegionId};
802
803 use super::*;
804 use crate::read::flat_projection::FlatProjectionMapper;
805 use crate::read::range::RowGroupIndex;
806 use crate::read::scan_region::{PredicateGroup, ScanInput};
807 use crate::read::scan_util::PartitionMetrics;
808 use crate::sst::file::{FileHandle, FileMeta};
809 use crate::sst::parquet::reader::ReaderMetrics;
810 use crate::test_util::memtable_util::metadata_with_primary_key;
811 use crate::test_util::new_noop_file_purger;
812 use crate::test_util::scheduler_util::SchedulerEnv;
813
814 async fn make_test_pruner(num_files: usize) -> (SchedulerEnv, Arc<Pruner>) {
815 make_test_pruner_with_retained_builders(num_files, false).await
816 }
817
818 async fn make_test_pruner_with_retained_builders(
819 num_files: usize,
820 retain_builders: bool,
821 ) -> (SchedulerEnv, Arc<Pruner>) {
822 let env = SchedulerEnv::new().await;
823 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
824 let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
825
826 let files: Vec<FileHandle> = (0..num_files)
827 .map(|_| {
828 let meta = FileMeta {
829 region_id: RegionId::new(123, 456),
830 file_id: FileId::random(),
831 time_range: (
832 Timestamp::new_millisecond(0),
833 Timestamp::new_millisecond(1000),
834 ),
835 num_row_groups: 1,
836 num_rows: 1024,
837 level: 0,
838 ..Default::default()
839 };
840 FileHandle::new(meta, new_noop_file_purger())
841 })
842 .collect();
843
844 let input = ScanInput::new(env.access_layer.clone(), mapper)
845 .with_files(files)
846 .with_append_mode(true);
847 let stream_ctx = Arc::new(StreamContext::unordered_scan_ctx(input));
848 let pruner = Arc::new(Pruner::new_with_options(
849 stream_ctx,
850 1,
851 PrunerOptions {
852 retain_builders,
853 ..Default::default()
854 },
855 ));
856 (env, pruner)
857 }
858
859 fn file_partition_range(file_index: usize) -> PartitionRange {
865 PartitionRange {
866 start: Timestamp::new_millisecond(0),
867 end: Timestamp::new_millisecond(1001),
868 num_rows: 1024,
869 identifier: file_index,
870 }
871 }
872
873 async fn make_test_pruner_with_predicate(
874 num_files: usize,
875 row_groups_per_file: u64,
876 predicate_exprs: &[Expr],
877 ) -> (SchedulerEnv, Arc<Pruner>) {
878 let env = SchedulerEnv::new().await;
879 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
880 let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
881 let predicate = PredicateGroup::new(&metadata, predicate_exprs).unwrap();
882
883 let files: Vec<FileHandle> = (0..num_files)
884 .map(|_| {
885 let meta = FileMeta {
886 region_id: RegionId::new(123, 456),
887 file_id: FileId::random(),
888 time_range: (
889 Timestamp::new_millisecond(0),
890 Timestamp::new_millisecond(1000),
891 ),
892 num_row_groups: row_groups_per_file,
893 num_rows: row_groups_per_file * 1024,
894 level: 0,
895 ..Default::default()
896 };
897 FileHandle::new(meta, new_noop_file_purger())
898 })
899 .collect();
900
901 let input = ScanInput::new(env.access_layer.clone(), mapper)
902 .with_files(files)
903 .with_predicate(predicate)
904 .with_append_mode(true);
905 let stream_ctx = Arc::new(StreamContext::unordered_scan_ctx(input));
906 let pruner = Arc::new(Pruner::new(stream_ctx, 1));
907 (env, pruner)
908 }
909
910 fn make_partition_metrics() -> PartitionMetrics {
911 let metrics_set = ExecutionPlanMetricsSet::new();
912 PartitionMetrics::new(
913 RegionId::new(123, 456),
914 0,
915 "test",
916 Instant::now(),
917 false,
918 &metrics_set,
919 )
920 }
921
922 #[test]
923 fn should_cache_builder_when_ranges_remain() {
924 let entry = FileBuilderEntry {
925 builder: None,
926 remaining_ranges: 3,
927 waiters: Vec::new(),
928 };
929 assert!(should_cache_builder(&entry));
930 }
931
932 #[test]
933 fn should_not_cache_builder_when_no_ranges_remain() {
934 let entry = FileBuilderEntry {
935 builder: None,
936 remaining_ranges: 0,
937 waiters: Vec::new(),
938 };
939 assert!(!should_cache_builder(&entry));
940 }
941
942 #[test]
943 fn should_not_cache_builder_when_already_cached() {
944 let entry = FileBuilderEntry {
945 builder: Some(Arc::new(FileRangeBuilder::default())),
946 remaining_ranges: 1,
947 waiters: Vec::new(),
948 };
949 assert!(!should_cache_builder(&entry));
950 }
951
952 #[test]
953 fn cache_builder_records_metrics() {
954 let mut entry = FileBuilderEntry {
955 builder: None,
956 remaining_ranges: 1,
957 waiters: Vec::new(),
958 };
959 let builder = Arc::new(FileRangeBuilder::default());
960 let mut reader_metrics = ReaderMetrics::default();
961
962 assert!(cache_builder_if_needed(
963 &mut entry,
964 &builder,
965 &mut reader_metrics,
966 ));
967 assert!(entry.builder.is_some());
968 assert_eq!(
969 reader_metrics.metadata_mem_size,
970 builder.memory_size() as isize
971 );
972 assert_eq!(reader_metrics.num_range_builders, 1);
973
974 if entry.builder.take().is_some() {
975 PRUNER_ACTIVE_BUILDERS.dec();
976 }
977 }
978
979 #[tokio::test]
980 async fn skip_file_range_decrements_and_clears_builder() {
981 let (_env, pruner) = make_test_pruner(1).await;
982
983 let ranges: Vec<PartitionRange> = (0..3).map(|_| file_partition_range(0)).collect();
985 pruner.add_partition_ranges(&ranges);
986 assert_eq!(pruner.test_remaining_ranges(0), 3);
987
988 {
990 let mut entry = pruner.inner.file_entries[0].lock().unwrap();
991 entry.builder = Some(Arc::new(FileRangeBuilder::default()));
992 PRUNER_ACTIVE_BUILDERS.inc();
993 }
994 assert!(pruner.test_has_builder(0));
995
996 let mut reader_metrics = ReaderMetrics::default();
998 for i in 0..3 {
999 let index = RowGroupIndex {
1000 index: 0,
1001 row_group_index: i as i64,
1002 };
1003 pruner.skip_file_range(index, &mut reader_metrics);
1004 }
1005
1006 assert_eq!(pruner.test_remaining_ranges(0), 0);
1007 assert!(!pruner.test_has_builder(0));
1008 }
1009
1010 #[tokio::test]
1011 async fn retained_builder_survives_after_last_range() {
1012 let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1013 pruner.add_partition_ranges(&[file_partition_range(0)]);
1014 {
1015 let mut entry = pruner.inner.file_entries[0].lock().unwrap();
1016 entry.builder = Some(Arc::new(FileRangeBuilder::default()));
1017 PRUNER_ACTIVE_BUILDERS.inc();
1018 }
1019
1020 let mut reader_metrics = ReaderMetrics::default();
1021 pruner.skip_file_range(
1022 RowGroupIndex {
1023 index: 0,
1024 row_group_index: 0,
1025 },
1026 &mut reader_metrics,
1027 );
1028
1029 assert_eq!(pruner.test_remaining_ranges(0), 0);
1030 assert!(pruner.test_has_builder(0));
1031
1032 let partition_metrics = make_partition_metrics();
1033 let mut reader_metrics = ReaderMetrics::default();
1034 let _builder = pruner
1035 .get_file_builder(
1036 0,
1037 PreFilterMode::SkipFields,
1038 &partition_metrics,
1039 &mut reader_metrics,
1040 )
1041 .await
1042 .unwrap();
1043 assert_eq!(reader_metrics.filter_metrics.pruner_cache_hit, 1);
1044 }
1045
1046 #[tokio::test]
1047 async fn add_partition_ranges_keeps_retained_builder() {
1048 let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1049 pruner.add_partition_ranges(&[file_partition_range(0)]);
1050 {
1051 let mut entry = pruner.inner.file_entries[0].lock().unwrap();
1052 entry.builder = Some(Arc::new(FileRangeBuilder::default()));
1053 PRUNER_ACTIVE_BUILDERS.inc();
1054 }
1055
1056 let mut reader_metrics = ReaderMetrics::default();
1057 pruner.skip_file_range(
1058 RowGroupIndex {
1059 index: 0,
1060 row_group_index: 0,
1061 },
1062 &mut reader_metrics,
1063 );
1064 assert!(pruner.test_has_builder(0));
1065
1066 pruner.add_partition_ranges(&[file_partition_range(0)]);
1067 assert!(pruner.test_has_builder(0));
1068 assert_eq!(pruner.test_remaining_ranges(0), 1);
1069 }
1070
1071 #[tokio::test]
1072 async fn retaining_mode_does_not_cache_after_skip_file_range_consumed_all() {
1073 let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1074
1075 let ranges = vec![file_partition_range(0)];
1077 pruner.add_partition_ranges(&ranges);
1078 assert_eq!(pruner.test_remaining_ranges(0), 1);
1079
1080 let mut reader_metrics = ReaderMetrics::default();
1084 let index = RowGroupIndex {
1085 index: 0,
1086 row_group_index: 0,
1087 };
1088 pruner.skip_file_range(index, &mut reader_metrics);
1089 assert_eq!(pruner.test_remaining_ranges(0), 0);
1090 assert!(!pruner.test_has_builder(0));
1091
1092 let entry = pruner.inner.file_entries[0].lock().unwrap();
1094 let should_cache = should_cache_builder(&entry);
1095 drop(entry);
1096
1097 assert!(!should_cache);
1098
1099 }
1103
1104 #[tokio::test]
1105 async fn worker_caches_when_ranges_remain() {
1106 let (_env, pruner) = make_test_pruner(1).await;
1107
1108 let ranges: Vec<PartitionRange> = (0..2).map(|_| file_partition_range(0)).collect();
1110 pruner.add_partition_ranges(&ranges);
1111 assert_eq!(pruner.test_remaining_ranges(0), 2);
1112
1113 let mut reader_metrics = ReaderMetrics::default();
1115 let index = RowGroupIndex {
1116 index: 0,
1117 row_group_index: 0,
1118 };
1119 pruner.skip_file_range(index, &mut reader_metrics);
1120 assert_eq!(pruner.test_remaining_ranges(0), 1);
1121
1122 let entry = pruner.inner.file_entries[0].lock().unwrap();
1124 assert!(should_cache_builder(&entry));
1125 }
1126
1127 #[tokio::test]
1133 async fn try_skip_manifest_pruned_file_range_multi_row_groups() {
1134 let predicate_exprs: Vec<Expr> =
1135 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1136 let (_env, pruner) = make_test_pruner_with_predicate(1, 3, &predicate_exprs).await;
1137
1138 let ranges = pruner.inner.stream_ctx.partition_ranges();
1139 assert_eq!(ranges.len(), 3);
1140 pruner.add_partition_ranges(&ranges);
1141 assert_eq!(pruner.test_remaining_ranges(0), 3);
1142
1143 let partition_pruner = Arc::new(PartitionPruner::new(pruner.clone(), &ranges));
1144 let partition_metrics = make_partition_metrics();
1145
1146 for rg in 0..3 {
1148 let index = RowGroupIndex {
1149 index: 0, row_group_index: rg,
1151 };
1152 let skipped =
1153 partition_pruner.try_skip_manifest_pruned_file_range(index, &partition_metrics);
1154 assert!(skipped, "row group {} should be skipped", rg);
1155 }
1156
1157 assert_eq!(pruner.test_remaining_ranges(0), 0);
1159 assert!(pruner.test_is_manifest_pruned(0));
1161 }
1162
1163 #[tokio::test]
1166 async fn try_skip_manifest_pruned_file_range_keeps_overlapping_file() {
1167 let predicate_exprs: Vec<Expr> =
1168 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(500), None)))];
1169 let (_env, pruner) = make_test_pruner_with_predicate(1, 2, &predicate_exprs).await;
1170
1171 let ranges = pruner.inner.stream_ctx.partition_ranges();
1172 assert_eq!(ranges.len(), 2);
1173 pruner.add_partition_ranges(&ranges);
1174 assert_eq!(pruner.test_remaining_ranges(0), 2);
1175
1176 let partition_pruner = Arc::new(PartitionPruner::new(pruner.clone(), &ranges));
1177 let partition_metrics = make_partition_metrics();
1178 let range_meta = &pruner.inner.stream_ctx.ranges[ranges[0].identifier];
1179 let index = range_meta.row_group_indices[0];
1180
1181 let skipped =
1182 partition_pruner.try_skip_manifest_pruned_file_range(index, &partition_metrics);
1183
1184 assert!(!skipped);
1185 assert_eq!(pruner.test_remaining_ranges(0), 2);
1186 assert!(!pruner.test_is_manifest_pruned(0));
1187 }
1188
1189 #[tokio::test]
1192 async fn add_partition_ranges_resets_manifest_pruned_flag() {
1193 let predicate_exprs: Vec<Expr> =
1194 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1195 let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1196
1197 let mut reader_metrics = ReaderMetrics::default();
1199 let marked = pruner
1200 .inner
1201 .try_mark_manifest_pruned(0, &mut reader_metrics);
1202 assert!(marked);
1203 assert!(pruner.test_is_manifest_pruned(0));
1204 assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1205
1206 let ranges = vec![file_partition_range(0)];
1208 pruner.add_partition_ranges(&ranges);
1209 assert!(!pruner.test_is_manifest_pruned(0));
1210 assert_eq!(pruner.test_remaining_ranges(0), 1);
1212 }
1213
1214 #[tokio::test]
1217 async fn prune_file_directly_manifest_pruned_returns_empty_builder() {
1218 let predicate_exprs: Vec<Expr> =
1219 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1220 let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1221
1222 assert!(!pruner.test_has_builder(0));
1224
1225 let mut reader_metrics = ReaderMetrics::default();
1226 let builder = pruner
1227 .prune_file_directly(0, PreFilterMode::SkipFields, &mut reader_metrics)
1228 .await
1229 .unwrap();
1230
1231 assert_eq!(
1233 builder.memory_size(),
1234 FileRangeBuilder::default().memory_size()
1235 );
1236 assert!(!pruner.test_has_builder(0));
1238 assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1240 }
1241
1242 #[tokio::test]
1245 async fn try_mark_manifest_pruned_only_counts_first_cas() {
1246 let predicate_exprs: Vec<Expr> =
1247 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1248 let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1249
1250 let mut reader_metrics = ReaderMetrics::default();
1252 let marked = pruner
1253 .inner
1254 .try_mark_manifest_pruned(0, &mut reader_metrics);
1255 assert!(marked);
1256 assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1257 assert!(pruner.test_is_manifest_pruned(0));
1258
1259 let mut reader_metrics2 = ReaderMetrics::default();
1261 let marked2 = pruner
1262 .inner
1263 .try_mark_manifest_pruned(0, &mut reader_metrics2);
1264 assert!(marked2);
1265 assert_eq!(reader_metrics2.filter_metrics.files_time_range_pruned, 0);
1266 }
1267}