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(crate) fn excluding_files(mut self, excluded: &HashSet<usize>) -> Self {
92 self.file_indices.retain(|index| !excluded.contains(index));
93 self
94 }
95
96 pub async fn build_file_ranges(
101 &self,
102 index: RowGroupIndex,
103 partition_metrics: &PartitionMetrics,
104 reader_metrics: &mut ReaderMetrics,
105 ) -> Result<SmallVec<[FileRange; 2]>> {
106 let file_index = index.index - self.pruner.inner.stream_ctx.input.num_memtables();
107 let pre_filter_mode = self.pre_filter_mode(file_index);
108
109 let ranges = self
111 .pruner
112 .build_file_ranges(index, pre_filter_mode, partition_metrics, reader_metrics)
113 .await?;
114
115 if let Some(pos) = self.file_indices.iter().position(|&idx| idx == file_index) {
117 let prev_pos = self.current_position.fetch_max(pos, Ordering::Relaxed);
118 if pos > prev_pos || prev_pos == 0 {
119 self.prefetch_upcoming_files(pos, partition_metrics);
120 }
121 }
122
123 Ok(ranges)
124 }
125
126 pub fn try_skip_manifest_pruned_file_range(
136 &self,
137 index: RowGroupIndex,
138 part_metrics: &PartitionMetrics,
139 ) -> bool {
140 let Some(file_index) = self.file_index(index) else {
141 return false;
142 };
143 let mut reader_metrics = ReaderMetrics::default();
144 let pruned = self
145 .pruner
146 .inner
147 .try_mark_manifest_pruned(file_index, &mut reader_metrics);
148 if pruned {
149 self.pruner.skip_file_range(index, &mut reader_metrics);
150 part_metrics.merge_reader_metrics(&reader_metrics, None);
151 }
152 pruned
153 }
154
155 fn prefetch_upcoming_files(&self, current_pos: usize, partition_metrics: &PartitionMetrics) {
157 let start = current_pos + 1;
158 let end = (start + PREFETCH_COUNT).min(self.file_indices.len());
159
160 for i in start..end {
161 let file_index = self.file_indices[i];
162 let pre_filter_mode = self.pre_filter_mode(file_index);
163 self.pruner.get_file_builder_background(
164 file_index,
165 pre_filter_mode,
166 Some(partition_metrics.clone()),
167 );
168 }
169 }
170
171 fn pre_filter_mode(&self, file_index: usize) -> PreFilterMode {
172 self.pre_filter_modes
173 .get(file_index)
174 .copied()
175 .unwrap_or(PreFilterMode::SkipFields)
176 }
177
178 fn file_index(&self, index: RowGroupIndex) -> Option<usize> {
179 self.pruner
180 .inner
181 .stream_ctx
182 .is_file_range_index(index)
183 .then(|| index.index - self.pruner.inner.stream_ctx.input.num_memtables())
184 }
185}
186
187#[derive(Debug, Clone, Copy)]
189pub struct PrunerOptions {
190 pub retain_builders: bool,
192 pub enable_predicate_prefilter: bool,
200}
201
202impl Default for PrunerOptions {
203 fn default() -> Self {
204 Self {
205 retain_builders: false,
206 enable_predicate_prefilter: true,
207 }
208 }
209}
210
211pub struct Pruner {
213 worker_senders: Vec<mpsc::Sender<PruneRequest>>,
215 inner: Arc<PrunerInner>,
216}
217
218struct PrunerInner {
219 num_workers: usize,
221 file_entries: Vec<Mutex<FileBuilderEntry>>,
223 stream_ctx: Arc<StreamContext>,
225 manifest_pruned_files: Vec<AtomicBool>,
231 retain_builders: bool,
233 enable_predicate_prefilter: bool,
235}
236
237impl Drop for PrunerInner {
238 fn drop(&mut self) {
239 let active_builders = self
240 .file_entries
241 .iter_mut()
242 .map(|entry| usize::from(entry.get_mut().is_ok_and(|entry| entry.builder.is_some())))
243 .sum::<usize>();
244 PRUNER_ACTIVE_BUILDERS.sub(active_builders as i64);
245 }
246}
247
248impl PrunerInner {
249 fn try_mark_manifest_pruned(
255 &self,
256 file_index: usize,
257 reader_metrics: &mut ReaderMetrics,
258 ) -> bool {
259 if self.manifest_pruned_files[file_index].load(Ordering::Relaxed) {
260 return true;
261 }
262 let file = &self.stream_ctx.input.files[file_index];
263 if !self.stream_ctx.input.can_manifest_prune_file(file) {
264 return false;
265 }
266 if self.manifest_pruned_files[file_index]
267 .compare_exchange(false, true, Ordering::Relaxed, Ordering::Relaxed)
268 .is_ok()
269 {
270 reader_metrics.filter_metrics.files_time_range_pruned += 1;
271 }
272 true
273 }
274}
275
276struct FileBuilderEntry {
278 builder: Option<Arc<FileRangeBuilder>>,
284 remaining_ranges: usize,
288 waiters: Vec<oneshot::Sender<Result<Arc<FileRangeBuilder>>>>,
290}
291
292struct PruneRequest {
294 file_index: usize,
296 pre_filter_mode: PreFilterMode,
298 response_tx: Option<oneshot::Sender<Result<Arc<FileRangeBuilder>>>>,
300 partition_metrics: Option<PartitionMetrics>,
302}
303
304impl Pruner {
305 pub fn new(stream_ctx: Arc<StreamContext>, num_workers: usize) -> Self {
310 Self::new_with_options(stream_ctx, num_workers, PrunerOptions::default())
311 }
312
313 pub fn new_with_options(
315 stream_ctx: Arc<StreamContext>,
316 num_workers: usize,
317 options: PrunerOptions,
318 ) -> Self {
319 let PrunerOptions {
320 retain_builders,
321 enable_predicate_prefilter,
322 } = options;
323 let num_files = stream_ctx.input.num_files();
324 let file_entries: Vec<_> = (0..num_files)
325 .map(|_| {
326 Mutex::new(FileBuilderEntry {
327 builder: None,
328 remaining_ranges: 0,
329 waiters: Vec::new(),
330 })
331 })
332 .collect();
333 let manifest_pruned_files: Vec<AtomicBool> =
334 (0..num_files).map(|_| AtomicBool::new(false)).collect();
335 let mut worker_senders = Vec::with_capacity(num_workers);
337 let mut receivers = Vec::with_capacity(num_workers);
338 for _ in 0..num_workers {
339 let (tx, rx) = mpsc::channel::<PruneRequest>(64);
340 worker_senders.push(tx);
341 receivers.push(rx);
342 }
343
344 let inner = Arc::new(PrunerInner {
345 num_workers,
346 file_entries,
347 stream_ctx,
348 manifest_pruned_files,
349 retain_builders,
350 enable_predicate_prefilter,
351 });
352
353 for (worker_id, rx) in receivers.into_iter().enumerate() {
355 let worker = Self::worker_loop(worker_id, rx, inner.clone());
356 if inner.stream_ctx.input.compaction {
357 common_runtime::spawn_compact(worker);
358 } else {
359 common_runtime::spawn_query(worker);
360 }
361 }
362
363 Self {
364 worker_senders,
365 inner,
366 }
367 }
368
369 pub fn add_partition_ranges(&self, partition_ranges: &[PartitionRange]) {
376 for pruned in &self.inner.manifest_pruned_files {
378 pruned.store(false, Ordering::Relaxed);
379 }
380
381 let num_memtables = self.inner.stream_ctx.input.num_memtables();
383 for part_range in partition_ranges {
384 let range_meta = &self.inner.stream_ctx.ranges[part_range.identifier];
385 for row_group_index in &range_meta.row_group_indices {
386 if self.inner.stream_ctx.is_file_range_index(*row_group_index) {
387 let file_index = row_group_index.index - num_memtables;
388 if file_index < self.inner.file_entries.len() {
389 let mut entry = self.inner.file_entries[file_index].lock().unwrap();
390 entry.remaining_ranges += 1;
391 }
392 }
393 }
394 }
395 }
396
397 pub async fn build_file_ranges(
403 &self,
404 index: RowGroupIndex,
405 pre_filter_mode: PreFilterMode,
406 partition_metrics: &PartitionMetrics,
407 reader_metrics: &mut ReaderMetrics,
408 ) -> Result<SmallVec<[FileRange; 2]>> {
409 let file_index = index.index - self.inner.stream_ctx.input.num_memtables();
410
411 let builder = self
413 .get_file_builder(
414 file_index,
415 pre_filter_mode,
416 partition_metrics,
417 reader_metrics,
418 )
419 .await?;
420
421 let mut ranges = SmallVec::new();
423 builder.build_ranges(index.row_group_index, &mut ranges);
424
425 self.decrement_and_maybe_clear(file_index, reader_metrics);
427
428 Ok(ranges)
429 }
430
431 pub fn skip_file_range(&self, index: RowGroupIndex, reader_metrics: &mut ReaderMetrics) {
437 if !self.inner.stream_ctx.is_file_range_index(index) {
438 return;
439 }
440 let file_index = index.index - self.inner.stream_ctx.input.num_memtables();
441 self.decrement_and_maybe_clear(file_index, reader_metrics);
442 }
443
444 async fn get_file_builder(
446 &self,
447 file_index: usize,
448 pre_filter_mode: PreFilterMode,
449 partition_metrics: &PartitionMetrics,
450 reader_metrics: &mut ReaderMetrics,
451 ) -> Result<Arc<FileRangeBuilder>> {
452 {
454 let entry = self.inner.file_entries[file_index].lock().unwrap();
455 if let Some(builder) = &entry.builder {
456 reader_metrics.filter_metrics.pruner_cache_hit += 1;
457 return Ok(builder.clone());
458 }
459 }
460
461 reader_metrics.filter_metrics.pruner_cache_miss += 1;
462 let prune_start = Instant::now();
463 let file = &self.inner.stream_ctx.input.files[file_index];
464 let file_id = file.file_id().file_id();
465 let worker_idx = self.get_worker_idx(file_id);
466
467 let (response_tx, response_rx) = oneshot::channel();
468 let request = PruneRequest {
469 file_index,
470 pre_filter_mode,
471 response_tx: Some(response_tx),
472 partition_metrics: Some(partition_metrics.clone()),
473 };
474
475 let result = if self.worker_senders[worker_idx].send(request).await.is_err() {
476 common_telemetry::warn!("Worker channel closed, falling back to direct pruning");
477 self.prune_file_directly(file_index, pre_filter_mode, reader_metrics)
479 .await
480 } else {
481 match response_rx.await {
483 Ok(result) => result,
484 Err(_) => {
485 common_telemetry::warn!(
486 "Response channel closed, falling back to direct pruning"
487 );
488 self.prune_file_directly(file_index, pre_filter_mode, reader_metrics)
490 .await
491 }
492 }
493 };
494 reader_metrics.filter_metrics.pruner_prune_cost += prune_start.elapsed();
495 result
496 }
497
498 pub fn get_file_builder_background(
500 &self,
501 file_index: usize,
502 pre_filter_mode: PreFilterMode,
503 partition_metrics: Option<PartitionMetrics>,
504 ) {
505 {
507 let entry = self.inner.file_entries[file_index].lock().unwrap();
508 if entry.builder.is_some() {
509 return;
510 }
511 }
512
513 let file = &self.inner.stream_ctx.input.files[file_index];
514 let file_id = file.file_id().file_id();
515 let worker_idx = self.get_worker_idx(file_id);
516
517 let request = PruneRequest {
518 file_index,
519 pre_filter_mode,
520 response_tx: None,
521 partition_metrics,
522 };
523
524 let _ = self.worker_senders[worker_idx].try_send(request);
526 }
527
528 pub fn predicate_prefilter_enabled(&self) -> bool {
531 self.inner.enable_predicate_prefilter
532 }
533
534 fn get_worker_idx(&self, file_id: FileId) -> usize {
535 let file_id_hash = Uuid::from(file_id).as_u128() as usize;
536 file_id_hash % self.inner.num_workers
537 }
538
539 async fn prune_file_directly(
542 &self,
543 file_index: usize,
544 pre_filter_mode: PreFilterMode,
545 reader_metrics: &mut ReaderMetrics,
546 ) -> Result<Arc<FileRangeBuilder>> {
547 if self
549 .inner
550 .try_mark_manifest_pruned(file_index, reader_metrics)
551 {
552 let arc_builder = Arc::new(FileRangeBuilder::default());
553 return Ok(arc_builder);
556 }
557
558 let file = &self.inner.stream_ctx.input.files[file_index];
559 let predicate = self.inner.stream_ctx.input.predicate_for_file(file);
560 let builder = self
561 .inner
562 .stream_ctx
563 .input
564 .prune_file_after_manifest_check(
565 file,
566 pre_filter_mode,
567 self.inner.enable_predicate_prefilter,
568 predicate,
569 reader_metrics,
570 )
571 .await?;
572
573 let arc_builder = Arc::new(builder);
574
575 {
578 let mut entry = self.inner.file_entries[file_index].lock().unwrap();
579 cache_builder_if_needed(&mut entry, &arc_builder, reader_metrics);
580 }
581
582 Ok(arc_builder)
583 }
584
585 fn decrement_and_maybe_clear(&self, file_index: usize, reader_metrics: &mut ReaderMetrics) {
587 let mut entry = self.inner.file_entries[file_index].lock().unwrap();
588 entry.remaining_ranges = entry.remaining_ranges.saturating_sub(1);
589
590 if !self.inner.retain_builders
591 && entry.remaining_ranges == 0
592 && let Some(builder) = entry.builder.take()
593 {
594 PRUNER_ACTIVE_BUILDERS.dec();
595 reader_metrics.metadata_mem_size -= builder.memory_size() as isize;
596 reader_metrics.num_range_builders -= 1;
597 }
598 }
599
600 async fn worker_loop(
602 worker_id: usize,
603 mut rx: mpsc::Receiver<PruneRequest>,
604 inner: Arc<PrunerInner>,
605 ) {
606 let mut worker_cache_hit = 0;
607 let mut worker_cache_miss = 0;
608 let mut pruned_files = Vec::new();
609
610 while let Some(request) = rx.recv().await {
611 let PruneRequest {
612 file_index,
613 pre_filter_mode,
614 response_tx,
615 partition_metrics,
616 } = request;
617
618 {
620 let entry = inner.file_entries[file_index].lock().unwrap();
621 if let Some(builder) = &entry.builder {
622 if let Some(response_tx) = response_tx {
624 let _ = response_tx.send(Ok(builder.clone()));
625 }
626 worker_cache_hit += 1;
627 continue;
628 }
629 }
630 worker_cache_miss += 1;
631
632 let file = &inner.stream_ctx.input.files[file_index];
633 pruned_files.push(file.file_id().file_id());
634 let explain_verbose = partition_metrics
635 .as_ref()
636 .map(|m| m.explain_verbose())
637 .unwrap_or(false);
638 let mut metrics = ReaderMetrics {
639 filter_metrics: new_filter_metrics(explain_verbose),
640 ..Default::default()
641 };
642
643 let result = if inner.try_mark_manifest_pruned(file_index, &mut metrics) {
645 Ok(FileRangeBuilder::default())
648 } else {
649 let predicate = inner.stream_ctx.input.predicate_for_file(file);
650 inner
651 .stream_ctx
652 .input
653 .prune_file_after_manifest_check(
654 file,
655 pre_filter_mode,
656 inner.enable_predicate_prefilter,
657 predicate,
658 &mut metrics,
659 )
660 .await
661 };
662
663 let mut entry = inner.file_entries[file_index].lock().unwrap();
665 match result {
666 Ok(builder) => {
667 let arc_builder = Arc::new(builder);
668 let is_background = response_tx.is_none();
669
670 let did_cache =
676 if inner.manifest_pruned_files[file_index].load(Ordering::Relaxed) {
677 false
678 } else {
679 cache_builder_if_needed(&mut entry, &arc_builder, &mut metrics)
680 };
681
682 for waiter in entry.waiters.drain(..) {
684 let _ = waiter.send(Ok(arc_builder.clone()));
685 }
686 if let Some(response_tx) = response_tx {
688 let _ = response_tx.send(Ok(arc_builder));
689 }
690
691 debug!(
692 "Pruner worker {} pruned file_index: {}, file: {:?}, metrics: {:?}",
693 worker_id,
694 file_index,
695 file.file_id(),
696 metrics
697 );
698
699 if (!is_background || did_cache)
704 && let Some(part_metrics) = &partition_metrics
705 {
706 let per_file_metrics = if part_metrics.explain_verbose() {
707 let file_id = file.file_id();
708 let mut map = HashMap::new();
709 map.insert(
710 file_id,
711 FileScanMetrics {
712 build_part_cost: metrics.build_cost,
713 ..Default::default()
714 },
715 );
716 Some(map)
717 } else {
718 None
719 };
720 part_metrics.merge_reader_metrics(&metrics, per_file_metrics.as_ref());
721 }
722 }
723 Err(e) => {
724 let arc_error = Arc::new(e);
725 for waiter in entry.waiters.drain(..) {
726 let _ = waiter.send(Err(arc_error.clone()).context(PruneFileSnafu));
727 }
728 if let Some(response_tx) = response_tx {
729 let _ = response_tx.send(Err(arc_error).context(PruneFileSnafu));
730 }
731 }
732 }
733 }
734
735 common_telemetry::debug!(
736 "Pruner worker {} finished, cache_hit: {}, cache_miss: {}, files: {:?}",
737 worker_id,
738 worker_cache_hit,
739 worker_cache_miss,
740 pruned_files,
741 );
742 }
743}
744
745#[cfg(test)]
746impl Pruner {
747 pub(crate) fn test_remaining_ranges(&self, file_index: usize) -> usize {
749 self.inner.file_entries[file_index]
750 .lock()
751 .unwrap()
752 .remaining_ranges
753 }
754
755 fn test_has_builder(&self, file_index: usize) -> bool {
757 self.inner.file_entries[file_index]
758 .lock()
759 .unwrap()
760 .builder
761 .is_some()
762 }
763
764 fn test_is_manifest_pruned(&self, file_index: usize) -> bool {
766 self.inner.manifest_pruned_files[file_index].load(Ordering::Relaxed)
767 }
768
769 #[allow(dead_code)]
771 fn test_clear_builder(&self, file_index: usize) {
772 let mut entry = self.inner.file_entries[file_index].lock().unwrap();
773 if entry.builder.take().is_some() {
774 PRUNER_ACTIVE_BUILDERS.dec();
775 }
776 }
777}
778
779fn should_cache_builder(entry: &FileBuilderEntry) -> bool {
781 entry.builder.is_none() && entry.remaining_ranges > 0
782}
783
784fn cache_builder_if_needed(
787 entry: &mut FileBuilderEntry,
788 builder: &Arc<FileRangeBuilder>,
789 reader_metrics: &mut ReaderMetrics,
790) -> bool {
791 if should_cache_builder(entry) {
792 reader_metrics.metadata_mem_size += builder.memory_size() as isize;
793 reader_metrics.num_range_builders += 1;
794 entry.builder = Some(builder.clone());
795 PRUNER_ACTIVE_BUILDERS.inc();
796 true
797 } else {
798 false
799 }
800}
801
802#[cfg(test)]
803mod tests {
804 use common_time::Timestamp;
805 use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet;
806 use datafusion_common::ScalarValue;
807 use datafusion_expr::{Expr, col, lit};
808 use store_api::region_engine::PartitionRange;
809 use store_api::storage::{FileId, RegionId};
810
811 use super::*;
812 use crate::read::flat_projection::FlatProjectionMapper;
813 use crate::read::range::RowGroupIndex;
814 use crate::read::scan_region::{PredicateGroup, ScanInput};
815 use crate::read::scan_util::PartitionMetrics;
816 use crate::sst::file::{FileHandle, FileMeta};
817 use crate::sst::parquet::reader::ReaderMetrics;
818 use crate::test_util::memtable_util::metadata_with_primary_key;
819 use crate::test_util::new_noop_file_purger;
820 use crate::test_util::scheduler_util::SchedulerEnv;
821
822 async fn make_test_pruner(num_files: usize) -> (SchedulerEnv, Arc<Pruner>) {
823 make_test_pruner_with_retained_builders(num_files, false).await
824 }
825
826 async fn make_test_pruner_with_retained_builders(
827 num_files: usize,
828 retain_builders: bool,
829 ) -> (SchedulerEnv, Arc<Pruner>) {
830 let env = SchedulerEnv::new().await;
831 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
832 let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
833
834 let files: Vec<FileHandle> = (0..num_files)
835 .map(|_| {
836 let meta = FileMeta {
837 region_id: RegionId::new(123, 456),
838 file_id: FileId::random(),
839 time_range: (
840 Timestamp::new_millisecond(0),
841 Timestamp::new_millisecond(1000),
842 ),
843 num_row_groups: 1,
844 num_rows: 1024,
845 level: 0,
846 ..Default::default()
847 };
848 FileHandle::new(meta, new_noop_file_purger())
849 })
850 .collect();
851
852 let input = ScanInput::builder(env.access_layer.clone(), mapper)
853 .with_files(files)
854 .with_append_mode(true)
855 .build();
856 let stream_ctx = Arc::new(StreamContext::unordered_scan_ctx(input));
857 let pruner = Arc::new(Pruner::new_with_options(
858 stream_ctx,
859 1,
860 PrunerOptions {
861 retain_builders,
862 ..Default::default()
863 },
864 ));
865 (env, pruner)
866 }
867
868 fn file_partition_range(file_index: usize) -> PartitionRange {
874 PartitionRange {
875 start: Timestamp::new_millisecond(0),
876 end: Timestamp::new_millisecond(1001),
877 num_rows: 1024,
878 identifier: file_index,
879 }
880 }
881
882 async fn make_test_pruner_with_predicate(
883 num_files: usize,
884 row_groups_per_file: u64,
885 predicate_exprs: &[Expr],
886 ) -> (SchedulerEnv, Arc<Pruner>) {
887 let env = SchedulerEnv::new().await;
888 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
889 let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
890 let predicate = PredicateGroup::new(&metadata, predicate_exprs).unwrap();
891
892 let files: Vec<FileHandle> = (0..num_files)
893 .map(|_| {
894 let meta = FileMeta {
895 region_id: RegionId::new(123, 456),
896 file_id: FileId::random(),
897 time_range: (
898 Timestamp::new_millisecond(0),
899 Timestamp::new_millisecond(1000),
900 ),
901 num_row_groups: row_groups_per_file,
902 num_rows: row_groups_per_file * 1024,
903 level: 0,
904 ..Default::default()
905 };
906 FileHandle::new(meta, new_noop_file_purger())
907 })
908 .collect();
909
910 let input = ScanInput::builder(env.access_layer.clone(), mapper)
911 .with_files(files)
912 .with_predicate(predicate)
913 .with_append_mode(true)
914 .build();
915 let stream_ctx = Arc::new(StreamContext::unordered_scan_ctx(input));
916 let pruner = Arc::new(Pruner::new(stream_ctx, 1));
917 (env, pruner)
918 }
919
920 fn make_partition_metrics() -> PartitionMetrics {
921 let metrics_set = ExecutionPlanMetricsSet::new();
922 PartitionMetrics::new(
923 RegionId::new(123, 456),
924 0,
925 "test",
926 Instant::now(),
927 false,
928 &metrics_set,
929 )
930 }
931
932 #[test]
933 fn should_cache_builder_when_ranges_remain() {
934 let entry = FileBuilderEntry {
935 builder: None,
936 remaining_ranges: 3,
937 waiters: Vec::new(),
938 };
939 assert!(should_cache_builder(&entry));
940 }
941
942 #[test]
943 fn should_not_cache_builder_when_no_ranges_remain() {
944 let entry = FileBuilderEntry {
945 builder: None,
946 remaining_ranges: 0,
947 waiters: Vec::new(),
948 };
949 assert!(!should_cache_builder(&entry));
950 }
951
952 #[test]
953 fn should_not_cache_builder_when_already_cached() {
954 let entry = FileBuilderEntry {
955 builder: Some(Arc::new(FileRangeBuilder::default())),
956 remaining_ranges: 1,
957 waiters: Vec::new(),
958 };
959 assert!(!should_cache_builder(&entry));
960 }
961
962 #[test]
963 fn cache_builder_records_metrics() {
964 let mut entry = FileBuilderEntry {
965 builder: None,
966 remaining_ranges: 1,
967 waiters: Vec::new(),
968 };
969 let builder = Arc::new(FileRangeBuilder::default());
970 let mut reader_metrics = ReaderMetrics::default();
971
972 assert!(cache_builder_if_needed(
973 &mut entry,
974 &builder,
975 &mut reader_metrics,
976 ));
977 assert!(entry.builder.is_some());
978 assert_eq!(
979 reader_metrics.metadata_mem_size,
980 builder.memory_size() as isize
981 );
982 assert_eq!(reader_metrics.num_range_builders, 1);
983
984 if entry.builder.take().is_some() {
985 PRUNER_ACTIVE_BUILDERS.dec();
986 }
987 }
988
989 #[tokio::test]
990 async fn skip_file_range_decrements_and_clears_builder() {
991 let (_env, pruner) = make_test_pruner(1).await;
992
993 let ranges: Vec<PartitionRange> = (0..3).map(|_| file_partition_range(0)).collect();
995 pruner.add_partition_ranges(&ranges);
996 assert_eq!(pruner.test_remaining_ranges(0), 3);
997
998 {
1000 let mut entry = pruner.inner.file_entries[0].lock().unwrap();
1001 entry.builder = Some(Arc::new(FileRangeBuilder::default()));
1002 PRUNER_ACTIVE_BUILDERS.inc();
1003 }
1004 assert!(pruner.test_has_builder(0));
1005
1006 let mut reader_metrics = ReaderMetrics::default();
1008 for i in 0..3 {
1009 let index = RowGroupIndex {
1010 index: 0,
1011 row_group_index: i as i64,
1012 };
1013 pruner.skip_file_range(index, &mut reader_metrics);
1014 }
1015
1016 assert_eq!(pruner.test_remaining_ranges(0), 0);
1017 assert!(!pruner.test_has_builder(0));
1018 }
1019
1020 #[tokio::test]
1021 async fn retained_builder_survives_after_last_range() {
1022 let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1023 pruner.add_partition_ranges(&[file_partition_range(0)]);
1024 {
1025 let mut entry = pruner.inner.file_entries[0].lock().unwrap();
1026 entry.builder = Some(Arc::new(FileRangeBuilder::default()));
1027 PRUNER_ACTIVE_BUILDERS.inc();
1028 }
1029
1030 let mut reader_metrics = ReaderMetrics::default();
1031 pruner.skip_file_range(
1032 RowGroupIndex {
1033 index: 0,
1034 row_group_index: 0,
1035 },
1036 &mut reader_metrics,
1037 );
1038
1039 assert_eq!(pruner.test_remaining_ranges(0), 0);
1040 assert!(pruner.test_has_builder(0));
1041
1042 let partition_metrics = make_partition_metrics();
1043 let mut reader_metrics = ReaderMetrics::default();
1044 let _builder = pruner
1045 .get_file_builder(
1046 0,
1047 PreFilterMode::SkipFields,
1048 &partition_metrics,
1049 &mut reader_metrics,
1050 )
1051 .await
1052 .unwrap();
1053 assert_eq!(reader_metrics.filter_metrics.pruner_cache_hit, 1);
1054 }
1055
1056 #[tokio::test]
1057 async fn add_partition_ranges_keeps_retained_builder() {
1058 let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1059 pruner.add_partition_ranges(&[file_partition_range(0)]);
1060 {
1061 let mut entry = pruner.inner.file_entries[0].lock().unwrap();
1062 entry.builder = Some(Arc::new(FileRangeBuilder::default()));
1063 PRUNER_ACTIVE_BUILDERS.inc();
1064 }
1065
1066 let mut reader_metrics = ReaderMetrics::default();
1067 pruner.skip_file_range(
1068 RowGroupIndex {
1069 index: 0,
1070 row_group_index: 0,
1071 },
1072 &mut reader_metrics,
1073 );
1074 assert!(pruner.test_has_builder(0));
1075
1076 pruner.add_partition_ranges(&[file_partition_range(0)]);
1077 assert!(pruner.test_has_builder(0));
1078 assert_eq!(pruner.test_remaining_ranges(0), 1);
1079 }
1080
1081 #[tokio::test]
1082 async fn retaining_mode_does_not_cache_after_skip_file_range_consumed_all() {
1083 let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1084
1085 let ranges = vec![file_partition_range(0)];
1087 pruner.add_partition_ranges(&ranges);
1088 assert_eq!(pruner.test_remaining_ranges(0), 1);
1089
1090 let mut reader_metrics = ReaderMetrics::default();
1094 let index = RowGroupIndex {
1095 index: 0,
1096 row_group_index: 0,
1097 };
1098 pruner.skip_file_range(index, &mut reader_metrics);
1099 assert_eq!(pruner.test_remaining_ranges(0), 0);
1100 assert!(!pruner.test_has_builder(0));
1101
1102 let entry = pruner.inner.file_entries[0].lock().unwrap();
1104 let should_cache = should_cache_builder(&entry);
1105 drop(entry);
1106
1107 assert!(!should_cache);
1108
1109 }
1113
1114 #[tokio::test]
1115 async fn worker_caches_when_ranges_remain() {
1116 let (_env, pruner) = make_test_pruner(1).await;
1117
1118 let ranges: Vec<PartitionRange> = (0..2).map(|_| file_partition_range(0)).collect();
1120 pruner.add_partition_ranges(&ranges);
1121 assert_eq!(pruner.test_remaining_ranges(0), 2);
1122
1123 let mut reader_metrics = ReaderMetrics::default();
1125 let index = RowGroupIndex {
1126 index: 0,
1127 row_group_index: 0,
1128 };
1129 pruner.skip_file_range(index, &mut reader_metrics);
1130 assert_eq!(pruner.test_remaining_ranges(0), 1);
1131
1132 let entry = pruner.inner.file_entries[0].lock().unwrap();
1134 assert!(should_cache_builder(&entry));
1135 }
1136
1137 #[tokio::test]
1143 async fn try_skip_manifest_pruned_file_range_multi_row_groups() {
1144 let predicate_exprs: Vec<Expr> =
1145 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1146 let (_env, pruner) = make_test_pruner_with_predicate(1, 3, &predicate_exprs).await;
1147
1148 let ranges = pruner.inner.stream_ctx.partition_ranges();
1149 assert_eq!(ranges.len(), 3);
1150 pruner.add_partition_ranges(&ranges);
1151 assert_eq!(pruner.test_remaining_ranges(0), 3);
1152
1153 let partition_pruner = Arc::new(PartitionPruner::new(pruner.clone(), &ranges));
1154 let partition_metrics = make_partition_metrics();
1155
1156 for rg in 0..3 {
1158 let index = RowGroupIndex {
1159 index: 0, row_group_index: rg,
1161 };
1162 let skipped =
1163 partition_pruner.try_skip_manifest_pruned_file_range(index, &partition_metrics);
1164 assert!(skipped, "row group {} should be skipped", rg);
1165 }
1166
1167 assert_eq!(pruner.test_remaining_ranges(0), 0);
1169 assert!(pruner.test_is_manifest_pruned(0));
1171 }
1172
1173 #[tokio::test]
1176 async fn try_skip_manifest_pruned_file_range_keeps_overlapping_file() {
1177 let predicate_exprs: Vec<Expr> =
1178 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(500), None)))];
1179 let (_env, pruner) = make_test_pruner_with_predicate(1, 2, &predicate_exprs).await;
1180
1181 let ranges = pruner.inner.stream_ctx.partition_ranges();
1182 assert_eq!(ranges.len(), 2);
1183 pruner.add_partition_ranges(&ranges);
1184 assert_eq!(pruner.test_remaining_ranges(0), 2);
1185
1186 let partition_pruner = Arc::new(PartitionPruner::new(pruner.clone(), &ranges));
1187 let partition_metrics = make_partition_metrics();
1188 let range_meta = &pruner.inner.stream_ctx.ranges[ranges[0].identifier];
1189 let index = range_meta.row_group_indices[0];
1190
1191 let skipped =
1192 partition_pruner.try_skip_manifest_pruned_file_range(index, &partition_metrics);
1193
1194 assert!(!skipped);
1195 assert_eq!(pruner.test_remaining_ranges(0), 2);
1196 assert!(!pruner.test_is_manifest_pruned(0));
1197 }
1198
1199 #[tokio::test]
1202 async fn add_partition_ranges_resets_manifest_pruned_flag() {
1203 let predicate_exprs: Vec<Expr> =
1204 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1205 let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1206
1207 let mut reader_metrics = ReaderMetrics::default();
1209 let marked = pruner
1210 .inner
1211 .try_mark_manifest_pruned(0, &mut reader_metrics);
1212 assert!(marked);
1213 assert!(pruner.test_is_manifest_pruned(0));
1214 assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1215
1216 let ranges = vec![file_partition_range(0)];
1218 pruner.add_partition_ranges(&ranges);
1219 assert!(!pruner.test_is_manifest_pruned(0));
1220 assert_eq!(pruner.test_remaining_ranges(0), 1);
1222 }
1223
1224 #[tokio::test]
1227 async fn prune_file_directly_manifest_pruned_returns_empty_builder() {
1228 let predicate_exprs: Vec<Expr> =
1229 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1230 let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1231
1232 assert!(!pruner.test_has_builder(0));
1234
1235 let mut reader_metrics = ReaderMetrics::default();
1236 let builder = pruner
1237 .prune_file_directly(0, PreFilterMode::SkipFields, &mut reader_metrics)
1238 .await
1239 .unwrap();
1240
1241 assert_eq!(
1243 builder.memory_size(),
1244 FileRangeBuilder::default().memory_size()
1245 );
1246 assert!(!pruner.test_has_builder(0));
1248 assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1250 }
1251
1252 #[tokio::test]
1255 async fn try_mark_manifest_pruned_only_counts_first_cas() {
1256 let predicate_exprs: Vec<Expr> =
1257 vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1258 let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1259
1260 let mut reader_metrics = ReaderMetrics::default();
1262 let marked = pruner
1263 .inner
1264 .try_mark_manifest_pruned(0, &mut reader_metrics);
1265 assert!(marked);
1266 assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1267 assert!(pruner.test_is_manifest_pruned(0));
1268
1269 let mut reader_metrics2 = ReaderMetrics::default();
1271 let marked2 = pruner
1272 .inner
1273 .try_mark_manifest_pruned(0, &mut reader_metrics2);
1274 assert!(marked2);
1275 assert_eq!(reader_metrics2.filter_metrics.files_time_range_pruned, 0);
1276 }
1277}