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::builder(env.access_layer.clone(), mapper)
845 .with_files(files)
846 .with_append_mode(true)
847 .build();
848 let stream_ctx = Arc::new(StreamContext::unordered_scan_ctx(input));
849 let pruner = Arc::new(Pruner::new_with_options(
850 stream_ctx,
851 1,
852 PrunerOptions {
853 retain_builders,
854 ..Default::default()
855 },
856 ));
857 (env, pruner)
858 }
859
860 fn file_partition_range(file_index: usize) -> PartitionRange {
866 PartitionRange {
867 start: Timestamp::new_millisecond(0),
868 end: Timestamp::new_millisecond(1001),
869 num_rows: 1024,
870 identifier: file_index,
871 }
872 }
873
874 async fn make_test_pruner_with_predicate(
875 num_files: usize,
876 row_groups_per_file: u64,
877 predicate_exprs: &[Expr],
878 ) -> (SchedulerEnv, Arc<Pruner>) {
879 let env = SchedulerEnv::new().await;
880 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
881 let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
882 let predicate = PredicateGroup::new(&metadata, predicate_exprs).unwrap();
883
884 let files: Vec<FileHandle> = (0..num_files)
885 .map(|_| {
886 let meta = FileMeta {
887 region_id: RegionId::new(123, 456),
888 file_id: FileId::random(),
889 time_range: (
890 Timestamp::new_millisecond(0),
891 Timestamp::new_millisecond(1000),
892 ),
893 num_row_groups: row_groups_per_file,
894 num_rows: row_groups_per_file * 1024,
895 level: 0,
896 ..Default::default()
897 };
898 FileHandle::new(meta, new_noop_file_purger())
899 })
900 .collect();
901
902 let input = ScanInput::builder(env.access_layer.clone(), mapper)
903 .with_files(files)
904 .with_predicate(predicate)
905 .with_append_mode(true)
906 .build();
907 let stream_ctx = Arc::new(StreamContext::unordered_scan_ctx(input));
908 let pruner = Arc::new(Pruner::new(stream_ctx, 1));
909 (env, pruner)
910 }
911
912 fn make_partition_metrics() -> PartitionMetrics {
913 let metrics_set = ExecutionPlanMetricsSet::new();
914 PartitionMetrics::new(
915 RegionId::new(123, 456),
916 0,
917 "test",
918 Instant::now(),
919 false,
920 &metrics_set,
921 )
922 }
923
924 #[test]
925 fn should_cache_builder_when_ranges_remain() {
926 let entry = FileBuilderEntry {
927 builder: None,
928 remaining_ranges: 3,
929 waiters: Vec::new(),
930 };
931 assert!(should_cache_builder(&entry));
932 }
933
934 #[test]
935 fn should_not_cache_builder_when_no_ranges_remain() {
936 let entry = FileBuilderEntry {
937 builder: None,
938 remaining_ranges: 0,
939 waiters: Vec::new(),
940 };
941 assert!(!should_cache_builder(&entry));
942 }
943
944 #[test]
945 fn should_not_cache_builder_when_already_cached() {
946 let entry = FileBuilderEntry {
947 builder: Some(Arc::new(FileRangeBuilder::default())),
948 remaining_ranges: 1,
949 waiters: Vec::new(),
950 };
951 assert!(!should_cache_builder(&entry));
952 }
953
954 #[test]
955 fn cache_builder_records_metrics() {
956 let mut entry = FileBuilderEntry {
957 builder: None,
958 remaining_ranges: 1,
959 waiters: Vec::new(),
960 };
961 let builder = Arc::new(FileRangeBuilder::default());
962 let mut reader_metrics = ReaderMetrics::default();
963
964 assert!(cache_builder_if_needed(
965 &mut entry,
966 &builder,
967 &mut reader_metrics,
968 ));
969 assert!(entry.builder.is_some());
970 assert_eq!(
971 reader_metrics.metadata_mem_size,
972 builder.memory_size() as isize
973 );
974 assert_eq!(reader_metrics.num_range_builders, 1);
975
976 if entry.builder.take().is_some() {
977 PRUNER_ACTIVE_BUILDERS.dec();
978 }
979 }
980
981 #[tokio::test]
982 async fn skip_file_range_decrements_and_clears_builder() {
983 let (_env, pruner) = make_test_pruner(1).await;
984
985 let ranges: Vec<PartitionRange> = (0..3).map(|_| file_partition_range(0)).collect();
987 pruner.add_partition_ranges(&ranges);
988 assert_eq!(pruner.test_remaining_ranges(0), 3);
989
990 {
992 let mut entry = pruner.inner.file_entries[0].lock().unwrap();
993 entry.builder = Some(Arc::new(FileRangeBuilder::default()));
994 PRUNER_ACTIVE_BUILDERS.inc();
995 }
996 assert!(pruner.test_has_builder(0));
997
998 let mut reader_metrics = ReaderMetrics::default();
1000 for i in 0..3 {
1001 let index = RowGroupIndex {
1002 index: 0,
1003 row_group_index: i as i64,
1004 };
1005 pruner.skip_file_range(index, &mut reader_metrics);
1006 }
1007
1008 assert_eq!(pruner.test_remaining_ranges(0), 0);
1009 assert!(!pruner.test_has_builder(0));
1010 }
1011
1012 #[tokio::test]
1013 async fn retained_builder_survives_after_last_range() {
1014 let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1015 pruner.add_partition_ranges(&[file_partition_range(0)]);
1016 {
1017 let mut entry = pruner.inner.file_entries[0].lock().unwrap();
1018 entry.builder = Some(Arc::new(FileRangeBuilder::default()));
1019 PRUNER_ACTIVE_BUILDERS.inc();
1020 }
1021
1022 let mut reader_metrics = ReaderMetrics::default();
1023 pruner.skip_file_range(
1024 RowGroupIndex {
1025 index: 0,
1026 row_group_index: 0,
1027 },
1028 &mut reader_metrics,
1029 );
1030
1031 assert_eq!(pruner.test_remaining_ranges(0), 0);
1032 assert!(pruner.test_has_builder(0));
1033
1034 let partition_metrics = make_partition_metrics();
1035 let mut reader_metrics = ReaderMetrics::default();
1036 let _builder = pruner
1037 .get_file_builder(
1038 0,
1039 PreFilterMode::SkipFields,
1040 &partition_metrics,
1041 &mut reader_metrics,
1042 )
1043 .await
1044 .unwrap();
1045 assert_eq!(reader_metrics.filter_metrics.pruner_cache_hit, 1);
1046 }
1047
1048 #[tokio::test]
1049 async fn add_partition_ranges_keeps_retained_builder() {
1050 let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1051 pruner.add_partition_ranges(&[file_partition_range(0)]);
1052 {
1053 let mut entry = pruner.inner.file_entries[0].lock().unwrap();
1054 entry.builder = Some(Arc::new(FileRangeBuilder::default()));
1055 PRUNER_ACTIVE_BUILDERS.inc();
1056 }
1057
1058 let mut reader_metrics = ReaderMetrics::default();
1059 pruner.skip_file_range(
1060 RowGroupIndex {
1061 index: 0,
1062 row_group_index: 0,
1063 },
1064 &mut reader_metrics,
1065 );
1066 assert!(pruner.test_has_builder(0));
1067
1068 pruner.add_partition_ranges(&[file_partition_range(0)]);
1069 assert!(pruner.test_has_builder(0));
1070 assert_eq!(pruner.test_remaining_ranges(0), 1);
1071 }
1072
1073 #[tokio::test]
1074 async fn retaining_mode_does_not_cache_after_skip_file_range_consumed_all() {
1075 let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1076
1077 let ranges = vec![file_partition_range(0)];
1079 pruner.add_partition_ranges(&ranges);
1080 assert_eq!(pruner.test_remaining_ranges(0), 1);
1081
1082 let mut reader_metrics = ReaderMetrics::default();
1086 let index = RowGroupIndex {
1087 index: 0,
1088 row_group_index: 0,
1089 };
1090 pruner.skip_file_range(index, &mut reader_metrics);
1091 assert_eq!(pruner.test_remaining_ranges(0), 0);
1092 assert!(!pruner.test_has_builder(0));
1093
1094 let entry = pruner.inner.file_entries[0].lock().unwrap();
1096 let should_cache = should_cache_builder(&entry);
1097 drop(entry);
1098
1099 assert!(!should_cache);
1100
1101 }
1105
1106 #[tokio::test]
1107 async fn worker_caches_when_ranges_remain() {
1108 let (_env, pruner) = make_test_pruner(1).await;
1109
1110 let ranges: Vec<PartitionRange> = (0..2).map(|_| file_partition_range(0)).collect();
1112 pruner.add_partition_ranges(&ranges);
1113 assert_eq!(pruner.test_remaining_ranges(0), 2);
1114
1115 let mut reader_metrics = ReaderMetrics::default();
1117 let index = RowGroupIndex {
1118 index: 0,
1119 row_group_index: 0,
1120 };
1121 pruner.skip_file_range(index, &mut reader_metrics);
1122 assert_eq!(pruner.test_remaining_ranges(0), 1);
1123
1124 let entry = pruner.inner.file_entries[0].lock().unwrap();
1126 assert!(should_cache_builder(&entry));
1127 }
1128
1129 #[tokio::test]
1135 async fn try_skip_manifest_pruned_file_range_multi_row_groups() {
1136 let predicate_exprs: Vec<Expr> =
1137 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1138 let (_env, pruner) = make_test_pruner_with_predicate(1, 3, &predicate_exprs).await;
1139
1140 let ranges = pruner.inner.stream_ctx.partition_ranges();
1141 assert_eq!(ranges.len(), 3);
1142 pruner.add_partition_ranges(&ranges);
1143 assert_eq!(pruner.test_remaining_ranges(0), 3);
1144
1145 let partition_pruner = Arc::new(PartitionPruner::new(pruner.clone(), &ranges));
1146 let partition_metrics = make_partition_metrics();
1147
1148 for rg in 0..3 {
1150 let index = RowGroupIndex {
1151 index: 0, row_group_index: rg,
1153 };
1154 let skipped =
1155 partition_pruner.try_skip_manifest_pruned_file_range(index, &partition_metrics);
1156 assert!(skipped, "row group {} should be skipped", rg);
1157 }
1158
1159 assert_eq!(pruner.test_remaining_ranges(0), 0);
1161 assert!(pruner.test_is_manifest_pruned(0));
1163 }
1164
1165 #[tokio::test]
1168 async fn try_skip_manifest_pruned_file_range_keeps_overlapping_file() {
1169 let predicate_exprs: Vec<Expr> =
1170 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(500), None)))];
1171 let (_env, pruner) = make_test_pruner_with_predicate(1, 2, &predicate_exprs).await;
1172
1173 let ranges = pruner.inner.stream_ctx.partition_ranges();
1174 assert_eq!(ranges.len(), 2);
1175 pruner.add_partition_ranges(&ranges);
1176 assert_eq!(pruner.test_remaining_ranges(0), 2);
1177
1178 let partition_pruner = Arc::new(PartitionPruner::new(pruner.clone(), &ranges));
1179 let partition_metrics = make_partition_metrics();
1180 let range_meta = &pruner.inner.stream_ctx.ranges[ranges[0].identifier];
1181 let index = range_meta.row_group_indices[0];
1182
1183 let skipped =
1184 partition_pruner.try_skip_manifest_pruned_file_range(index, &partition_metrics);
1185
1186 assert!(!skipped);
1187 assert_eq!(pruner.test_remaining_ranges(0), 2);
1188 assert!(!pruner.test_is_manifest_pruned(0));
1189 }
1190
1191 #[tokio::test]
1194 async fn add_partition_ranges_resets_manifest_pruned_flag() {
1195 let predicate_exprs: Vec<Expr> =
1196 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1197 let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1198
1199 let mut reader_metrics = ReaderMetrics::default();
1201 let marked = pruner
1202 .inner
1203 .try_mark_manifest_pruned(0, &mut reader_metrics);
1204 assert!(marked);
1205 assert!(pruner.test_is_manifest_pruned(0));
1206 assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1207
1208 let ranges = vec![file_partition_range(0)];
1210 pruner.add_partition_ranges(&ranges);
1211 assert!(!pruner.test_is_manifest_pruned(0));
1212 assert_eq!(pruner.test_remaining_ranges(0), 1);
1214 }
1215
1216 #[tokio::test]
1219 async fn prune_file_directly_manifest_pruned_returns_empty_builder() {
1220 let predicate_exprs: Vec<Expr> =
1221 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1222 let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1223
1224 assert!(!pruner.test_has_builder(0));
1226
1227 let mut reader_metrics = ReaderMetrics::default();
1228 let builder = pruner
1229 .prune_file_directly(0, PreFilterMode::SkipFields, &mut reader_metrics)
1230 .await
1231 .unwrap();
1232
1233 assert_eq!(
1235 builder.memory_size(),
1236 FileRangeBuilder::default().memory_size()
1237 );
1238 assert!(!pruner.test_has_builder(0));
1240 assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1242 }
1243
1244 #[tokio::test]
1247 async fn try_mark_manifest_pruned_only_counts_first_cas() {
1248 let predicate_exprs: Vec<Expr> =
1249 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1250 let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1251
1252 let mut reader_metrics = ReaderMetrics::default();
1254 let marked = pruner
1255 .inner
1256 .try_mark_manifest_pruned(0, &mut reader_metrics);
1257 assert!(marked);
1258 assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1259 assert!(pruner.test_is_manifest_pruned(0));
1260
1261 let mut reader_metrics2 = ReaderMetrics::default();
1263 let marked2 = pruner
1264 .inner
1265 .try_mark_manifest_pruned(0, &mut reader_metrics2);
1266 assert!(marked2);
1267 assert_eq!(reader_metrics2.filter_metrics.files_time_range_pruned, 0);
1268 }
1269}