1use std::collections::{BTreeMap, HashSet};
18use std::fmt;
19use std::num::NonZeroU64;
20use std::sync::{Arc, OnceLock};
21use std::time::Instant;
22
23use api::v1::SemanticType;
24use common_error::ext::BoxedError;
25use common_recordbatch::SendableRecordBatchStream;
26use common_recordbatch::adapter::RegionQueryStatCounters;
27use common_recordbatch::filter::SimpleFilterEvaluator;
28use common_telemetry::tracing::Instrument;
29use common_telemetry::{debug, error, tracing, warn};
30use common_time::range::TimestampRange;
31use datafusion::execution::memory_pool::{MemoryPool, UnboundedMemoryPool};
32use datafusion::physical_plan::expressions::DynamicFilterPhysicalExpr;
33use datafusion_common::pruning::PruningStatistics;
34use datafusion_common::{Column, ScalarValue};
35use datafusion_expr::Expr;
36use datafusion_expr::utils::expr_to_columns;
37use datatypes::arrow::array::{ArrayRef, BooleanArray, UInt64Array};
38use datatypes::extension::json::is_json2_extension_type;
39use datatypes::prelude::ConcreteDataType;
40use datatypes::types::json_type::JsonNativeType;
41use datatypes::value::timestamp_to_scalar_value;
42use futures::StreamExt;
43use partition::expr::PartitionExpr;
44use smallvec::SmallVec;
45use snafu::{OptionExt, ResultExt, ensure};
46use store_api::metadata::{RegionMetadata, RegionMetadataRef};
47use store_api::region_engine::{PartitionRange, RegionScannerRef};
48use store_api::storage::{
49 ColumnId, RegionId, ScanRequest, SequenceNumber, SequenceRange, TimeSeriesDistribution,
50 TimeSeriesRowSelector,
51};
52use table::predicate::{Predicate, build_time_range_predicate, extract_time_range_from_expr};
53use tokio::sync::{Semaphore, mpsc};
54use tokio_stream::wrappers::ReceiverStream;
55
56use crate::access_layer::AccessLayerRef;
57use crate::cache::CacheStrategy;
58use crate::config::DEFAULT_MAX_CONCURRENT_SCAN_FILES;
59use crate::error::{
60 InvalidPartitionExprSnafu, InvalidRequestSnafu, RegionSequenceDomainBrokenSnafu, Result,
61 SequenceRangeUnsupportedSnafu,
62};
63#[cfg(feature = "enterprise")]
64use crate::extension::{BoxedExtensionRange, BoxedExtensionRangeProvider};
65use crate::memtable::{MemtableRange, RangesOptions};
66use crate::metrics::READ_SST_COUNT;
67use crate::read::compat::{self, FlatCompatBatch};
68use crate::read::flat_projection::FlatProjectionMapper;
69use crate::read::range::{FileRangeBuilder, MemRangeBuilder, RangeMeta, RowGroupIndex};
70use crate::read::range_cache::{ScanRequestFingerprint, implied_time_range_from_exprs};
71use crate::read::read_columns::ReadColumns;
72use crate::read::seq_scan::SeqScan;
73use crate::read::series_scan::SeriesScan;
74use crate::read::stream::ScanBatchStream;
75use crate::read::unordered_scan::UnorderedScan;
76use crate::read::{BoxedRecordBatchStream, RecordBatch};
77use crate::region::options::MergeMode;
78use crate::region::version::VersionRef;
79use crate::series_index::SeriesIndexReadContext;
80use crate::sst::file::FileHandle;
81use crate::sst::index::bloom_filter::applier::{
82 BloomFilterIndexApplierBuilder, BloomFilterIndexApplierRef,
83};
84use crate::sst::index::fulltext_index::applier::FulltextIndexApplierRef;
85use crate::sst::index::fulltext_index::applier::builder::FulltextIndexApplierBuilder;
86use crate::sst::index::inverted_index::applier::InvertedIndexApplierRef;
87use crate::sst::index::inverted_index::applier::builder::InvertedIndexApplierBuilder;
88#[cfg(feature = "vector_index")]
89use crate::sst::index::vector_index::applier::{VectorIndexApplier, VectorIndexApplierRef};
90use crate::sst::parquet::Json2RewriteTargets;
91use crate::sst::parquet::file_range::PreFilterMode;
92use crate::sst::parquet::reader::ReaderMetrics;
93use crate::sst::primary_key::PrimaryKeyRangeMapper;
94
95#[cfg(feature = "vector_index")]
96const VECTOR_INDEX_OVERFETCH_MULTIPLIER: usize = 2;
97
98pub(crate) enum Scanner {
100 Seq(SeqScan),
102 Unordered(UnorderedScan),
104 Series(SeriesScan),
106}
107
108impl Scanner {
109 #[tracing::instrument(level = tracing::Level::DEBUG, skip_all)]
111 pub(crate) async fn scan(&self) -> Result<SendableRecordBatchStream, BoxedError> {
112 match self {
113 Scanner::Seq(seq_scan) => seq_scan.build_stream(),
114 Scanner::Unordered(unordered_scan) => unordered_scan.build_stream().await,
115 Scanner::Series(series_scan) => series_scan.build_stream().await,
116 }
117 }
118
119 pub(crate) fn scan_batch(&self) -> Result<ScanBatchStream> {
121 match self {
122 Scanner::Seq(x) => x.scan_all_partitions(),
123 Scanner::Unordered(x) => x.scan_all_partitions(),
124 Scanner::Series(x) => x.scan_all_partitions(),
125 }
126 }
127}
128
129#[cfg(test)]
130impl Scanner {
131 pub(crate) fn num_files(&self) -> usize {
133 match self {
134 Scanner::Seq(seq_scan) => seq_scan.input().num_files(),
135 Scanner::Unordered(unordered_scan) => unordered_scan.input().num_files(),
136 Scanner::Series(series_scan) => series_scan.input().num_files(),
137 }
138 }
139
140 pub(crate) fn num_memtables(&self) -> usize {
142 match self {
143 Scanner::Seq(seq_scan) => seq_scan.input().num_memtables(),
144 Scanner::Unordered(unordered_scan) => unordered_scan.input().num_memtables(),
145 Scanner::Series(series_scan) => series_scan.input().num_memtables(),
146 }
147 }
148
149 pub(crate) fn file_ids(&self) -> Vec<crate::sst::file::RegionFileId> {
151 match self {
152 Scanner::Seq(seq_scan) => seq_scan.input().file_ids(),
153 Scanner::Unordered(unordered_scan) => unordered_scan.input().file_ids(),
154 Scanner::Series(series_scan) => series_scan.input().file_ids(),
155 }
156 }
157
158 pub(crate) fn index_ids(&self) -> Vec<crate::sst::file::RegionIndexId> {
159 match self {
160 Scanner::Seq(seq_scan) => seq_scan.input().index_ids(),
161 Scanner::Unordered(unordered_scan) => unordered_scan.input().index_ids(),
162 Scanner::Series(series_scan) => series_scan.input().index_ids(),
163 }
164 }
165
166 pub(crate) fn snapshot_sequence(&self) -> Option<SequenceNumber> {
167 match self {
168 Scanner::Seq(seq_scan) => seq_scan.input().snapshot_sequence,
169 Scanner::Unordered(unordered_scan) => unordered_scan.input().snapshot_sequence,
170 Scanner::Series(series_scan) => series_scan.input().snapshot_sequence,
171 }
172 }
173
174 pub(crate) fn set_target_partitions(&mut self, target_partitions: usize) {
176 use store_api::region_engine::{PrepareRequest, RegionScanner};
177
178 let request = PrepareRequest::default().with_target_partitions(target_partitions);
179 match self {
180 Scanner::Seq(seq_scan) => seq_scan.prepare(request).unwrap(),
181 Scanner::Unordered(unordered_scan) => unordered_scan.prepare(request).unwrap(),
182 Scanner::Series(series_scan) => series_scan.prepare(request).unwrap(),
183 }
184 }
185}
186
187#[cfg_attr(doc, aquamarine::aquamarine)]
188pub(crate) struct ScanRegion {
238 version: VersionRef,
240 series_index: Option<SeriesIndexReadContext>,
242 access_layer: AccessLayerRef,
244 request: ScanRequest,
246 cache_strategy: CacheStrategy,
248 max_concurrent_scan_files: usize,
250 scan_memory_pool: Arc<dyn MemoryPool>,
252 experimental_series_scan_v2: bool,
254 ignore_range_index: bool,
256 ignore_inverted_index: bool,
258 ignore_fulltext_index: bool,
260 ignore_bloom_filter: bool,
262 start_time: Option<Instant>,
264 filter_deleted: bool,
267 exact_selection: Option<(Vec<FileHandle>, Option<SequenceRange>)>,
269 query_stat_counters: Option<RegionQueryStatCounters>,
271 #[cfg(feature = "enterprise")]
272 extension_range_provider: Option<BoxedExtensionRangeProvider>,
273}
274
275impl ScanRegion {
276 pub(crate) fn new(
278 version: VersionRef,
279 access_layer: AccessLayerRef,
280 request: ScanRequest,
281 cache_strategy: CacheStrategy,
282 ) -> ScanRegion {
283 ScanRegion {
284 version,
285 series_index: None,
286 access_layer,
287 request,
288 cache_strategy,
289 max_concurrent_scan_files: DEFAULT_MAX_CONCURRENT_SCAN_FILES,
290 scan_memory_pool: Arc::new(UnboundedMemoryPool::default()),
291 experimental_series_scan_v2: false,
292 ignore_range_index: false,
293 ignore_inverted_index: false,
294 ignore_fulltext_index: false,
295 ignore_bloom_filter: false,
296 start_time: None,
297 filter_deleted: true,
298 exact_selection: None,
299 query_stat_counters: None,
300 #[cfg(feature = "enterprise")]
301 extension_range_provider: None,
302 }
303 }
304
305 #[must_use]
307 pub(crate) fn with_series_index(
308 mut self,
309 series_index: Option<SeriesIndexReadContext>,
310 ) -> Self {
311 self.series_index = series_index;
312 self
313 }
314
315 #[must_use]
317 pub(crate) fn with_query_stat_counters(mut self, counters: RegionQueryStatCounters) -> Self {
318 self.query_stat_counters = Some(counters);
319 self
320 }
321
322 #[must_use]
324 pub(crate) fn with_max_concurrent_scan_files(
325 mut self,
326 max_concurrent_scan_files: usize,
327 ) -> Self {
328 self.max_concurrent_scan_files = max_concurrent_scan_files;
329 self
330 }
331
332 #[must_use]
334 pub(crate) fn with_scan_memory_pool(mut self, scan_memory_pool: Arc<dyn MemoryPool>) -> Self {
335 self.scan_memory_pool = scan_memory_pool;
336 self
337 }
338
339 #[must_use]
341 pub(crate) fn with_experimental_series_scan_v2(mut self, enabled: bool) -> Self {
342 self.experimental_series_scan_v2 = enabled;
343 self
344 }
345
346 #[must_use]
348 pub(crate) fn with_ignore_range_index(mut self, ignore: bool) -> Self {
349 self.ignore_range_index = ignore;
350 self
351 }
352
353 #[must_use]
355 pub(crate) fn with_ignore_inverted_index(mut self, ignore: bool) -> Self {
356 self.ignore_inverted_index = ignore;
357 self
358 }
359
360 #[must_use]
362 pub(crate) fn with_ignore_fulltext_index(mut self, ignore: bool) -> Self {
363 self.ignore_fulltext_index = ignore;
364 self
365 }
366
367 #[must_use]
369 pub(crate) fn with_ignore_bloom_filter(mut self, ignore: bool) -> Self {
370 self.ignore_bloom_filter = ignore;
371 self
372 }
373
374 #[must_use]
375 pub(crate) fn with_start_time(mut self, now: Instant) -> Self {
376 self.start_time = Some(now);
377 self
378 }
379
380 pub(crate) fn set_filter_deleted(&mut self, filter_deleted: bool) {
381 self.filter_deleted = filter_deleted;
382 }
383
384 pub(crate) fn with_exact_selection(
385 mut self,
386 exact_selection: (Vec<FileHandle>, Option<SequenceRange>),
387 ) -> Self {
388 self.exact_selection = Some(exact_selection);
389 self
390 }
391
392 #[cfg(feature = "enterprise")]
393 pub(crate) fn set_extension_range_provider(
394 &mut self,
395 extension_range_provider: BoxedExtensionRangeProvider,
396 ) {
397 self.extension_range_provider = Some(extension_range_provider);
398 }
399
400 #[tracing::instrument(skip_all, fields(region_id = %self.region_id()))]
402 pub(crate) async fn scanner(self) -> Result<Scanner> {
403 if self.use_series_scan() {
404 self.series_scan().await.map(Scanner::Series)
405 } else if self.use_unordered_scan() {
406 self.unordered_scan().await.map(Scanner::Unordered)
409 } else {
410 self.seq_scan().await.map(Scanner::Seq)
411 }
412 }
413
414 #[tracing::instrument(
416 level = tracing::Level::DEBUG,
417 skip_all,
418 fields(region_id = %self.region_id())
419 )]
420 pub(crate) async fn region_scanner(self) -> Result<RegionScannerRef> {
421 if self.use_series_scan() {
422 self.series_scan()
423 .await
424 .map(|scanner| Box::new(scanner) as _)
425 } else if self.use_unordered_scan() {
426 self.unordered_scan()
427 .await
428 .map(|scanner| Box::new(scanner) as _)
429 } else {
430 self.seq_scan().await.map(|scanner| Box::new(scanner) as _)
431 }
432 }
433
434 #[tracing::instrument(skip_all, fields(region_id = %self.region_id()))]
436 pub(crate) async fn seq_scan(self) -> Result<SeqScan> {
437 let input = self.scan_input().await?;
438 Ok(SeqScan::new(input))
439 }
440
441 #[tracing::instrument(skip_all, fields(region_id = %self.region_id()))]
443 pub(crate) async fn unordered_scan(self) -> Result<UnorderedScan> {
444 let input = self.scan_input().await?;
445 Ok(UnorderedScan::new(input))
446 }
447
448 #[tracing::instrument(skip_all, fields(region_id = %self.region_id()))]
450 pub(crate) async fn series_scan(self) -> Result<SeriesScan> {
451 let experimental_series_scan_v2 = self.experimental_series_scan_v2;
452 let input = self.scan_input().await?;
453 Ok(SeriesScan::new(input, experimental_series_scan_v2))
454 }
455
456 fn use_unordered_scan(&self) -> bool {
458 self.version.options.append_mode
465 && self.request.series_row_selector.is_none()
466 && (self.request.distribution.is_none()
467 || self.request.distribution == Some(TimeSeriesDistribution::TimeWindowed))
468 }
469
470 fn use_series_scan(&self) -> bool {
472 self.request.distribution == Some(TimeSeriesDistribution::PerSeries)
473 }
474
475 #[tracing::instrument(skip_all, fields(region_id = %self.region_id()))]
477 async fn scan_input(mut self) -> Result<ScanInput> {
478 let metadata = &self.version.metadata;
479 let sst_min_sequence = self.request.sst_min_sequence.and_then(NonZeroU64::new);
480 let time_range = self.build_time_range_predicate();
481 let predicate = PredicateGroup::new(metadata, &self.request.filters)?;
482
483 let read_col_ids =
484 self.build_read_col_ids(self.request.projection.as_deref(), &predicate)?;
485 let read_cols = self.build_read_columns(&read_col_ids)?;
486
487 let projection = self
489 .request
490 .projection
491 .clone()
492 .unwrap_or_else(|| (0..metadata.column_metadatas.len()).collect());
493 let mapper = FlatProjectionMapper::new_with_read_columns(metadata, projection, read_cols)?;
494 let mapper = if self.request.preserve_pk_dictionary_encoding {
495 mapper.with_pk_dictionary_encoding()
496 } else {
497 mapper
498 };
499
500 let (files, sequence_range) =
501 if let Some((files, sequence_range)) = self.exact_selection.take() {
502 (files, sequence_range)
503 } else {
504 exact_sequence_range(&self.request, &self.version)?
505 };
506 if sst_min_sequence.is_some() && sequence_range.is_some() {
507 return SequenceRangeUnsupportedSnafu {
508 region_id: self.region_id(),
509 min_seq: self.request.memtable_min_sequence.unwrap_or_default(),
510 max_seq: self.request.memtable_max_sequence.unwrap_or_default(),
511 reason:
512 "sst_min_sequence pruning hint is incompatible with exact sequence-range reads"
513 .to_string(),
514 }
515 .fail();
516 }
517 let memtables = self.version.memtables.list_memtables();
518 let mut mem_range_builders = Vec::new();
520 let filter_mode = pre_filter_mode(
521 self.version.options.append_mode,
522 self.version.options.merge_mode(),
523 );
524
525 for m in memtables {
526 let Some((start, end)) = m.stats().time_range() else {
528 continue;
529 };
530 let memtable_range = TimestampRange::new_inclusive(Some(start), Some(end));
532 if !memtable_range.intersects(&time_range) {
533 continue;
534 }
535 let ranges_in_memtable = m.ranges(
536 Some(&read_col_ids),
537 RangesOptions::default()
538 .with_predicate(predicate.clone())
539 .with_sequence(SequenceRange::new(
540 self.request.memtable_min_sequence,
541 self.request.memtable_max_sequence,
542 ))
543 .with_pre_filter_mode(filter_mode),
544 )?;
545 mem_range_builders.extend(ranges_in_memtable.ranges.into_values().map(|v| {
546 let stats = v.stats().clone();
547 MemRangeBuilder::new(v, stats)
548 }));
549 }
550
551 let region_id = self.region_id();
552 debug!(
553 "Scan region {}, request: {:?}, time range: {:?}, memtables: {}, ssts_to_read: {}, append_mode: {}",
554 region_id,
555 self.request,
556 time_range,
557 mem_range_builders.len(),
558 files.len(),
559 self.version.options.append_mode,
560 );
561
562 let (non_field_filters, field_filters) = self.partition_by_field_filters();
563 let inverted_index_appliers = [
564 self.build_invereted_index_applier(&non_field_filters),
565 self.build_invereted_index_applier(&field_filters),
566 ];
567 let bloom_filter_appliers = [
568 self.build_bloom_filter_applier(&non_field_filters),
569 self.build_bloom_filter_applier(&field_filters),
570 ];
571 let fulltext_index_appliers = [
572 self.build_fulltext_index_applier(&non_field_filters),
573 self.build_fulltext_index_applier(&field_filters),
574 ];
575 #[cfg(feature = "vector_index")]
576 let vector_index_applier = self.build_vector_index_applier();
577 #[cfg(feature = "vector_index")]
578 let vector_index_k = self.request.vector_search.as_ref().map(|search| {
579 if self.request.filters.is_empty() {
580 search.k
581 } else {
582 search.k.saturating_mul(VECTOR_INDEX_OVERFETCH_MULTIPLIER)
583 }
584 });
585
586 let input = ScanInput::builder(self.access_layer, mapper)
587 .with_series_index(self.series_index)
588 .with_ignore_range_index(self.ignore_range_index)
589 .with_time_range(Some(time_range))
590 .with_predicate(predicate)
591 .with_memtables(mem_range_builders)
592 .with_files(files)
593 .with_primary_key_mapper(self.version.ssts.primary_key_mapper())
594 .with_cache(self.cache_strategy)
595 .with_inverted_index_appliers(inverted_index_appliers)
596 .with_bloom_filter_index_appliers(bloom_filter_appliers)
597 .with_fulltext_index_appliers(fulltext_index_appliers)
598 .with_max_concurrent_scan_files(self.max_concurrent_scan_files)
599 .with_scan_memory_pool(self.scan_memory_pool)
600 .with_start_time(self.start_time)
601 .with_append_mode(self.version.options.append_mode)
602 .with_filter_deleted(self.filter_deleted)
603 .with_merge_mode(self.version.options.merge_mode())
604 .with_series_row_selector(self.request.series_row_selector)
605 .with_distribution(self.request.distribution)
606 .with_explain_flat_format(
607 self.version.options.sst_format == Some(crate::sst::FormatType::Flat),
608 )
609 .with_snapshot_sequence(
610 self.request
611 .snapshot_on_scan
612 .then_some(self.request.memtable_max_sequence)
613 .flatten(),
614 )
615 .with_sequence_range(sequence_range)
616 .with_query_stat_counters(self.query_stat_counters);
617 #[cfg(feature = "vector_index")]
618 let input = input
619 .with_vector_index_applier(vector_index_applier)
620 .with_vector_index_k(vector_index_k);
621
622 #[cfg(feature = "enterprise")]
623 let input = if !self.request.skip_sst_files
624 && let Some(provider) = self.extension_range_provider
625 {
626 if sequence_range.is_some() {
627 return SequenceRangeUnsupportedSnafu {
633 region_id,
634 min_seq: self.request.memtable_min_sequence.unwrap_or_default(),
635 max_seq: self.request.memtable_max_sequence.unwrap_or_default(),
636 reason:
637 "exact sequence-range reads are unsupported when an extension range provider is present"
638 .to_string(),
639 }
640 .fail();
641 }
642 let ranges = provider
643 .find_extension_ranges(self.version.flushed_sequence, time_range, &self.request)
644 .await?;
645 debug!("Find extension ranges: {ranges:?}");
646 input.with_extension_ranges(ranges)
647 } else {
648 input
649 };
650 Ok(input.build())
651 }
652
653 fn build_read_col_ids(
656 &self,
657 projection: Option<&[usize]>,
658 predicate: &PredicateGroup,
659 ) -> Result<Vec<ColumnId>> {
660 let metadata = &self.version.metadata;
661 let Some(projection) = projection else {
662 return Ok(metadata
663 .column_metadatas
664 .iter()
665 .map(|col| col.column_id)
666 .collect());
667 };
668
669 let mut read_col_ids = Vec::new();
670 let mut seen = HashSet::new();
671 for idx in projection {
672 let col_id = metadata
673 .column_metadatas
674 .get(*idx)
675 .with_context(|| InvalidRequestSnafu {
676 region_id: metadata.region_id,
677 reason: format!("projection index {} is out of bounds", idx),
678 })?
679 .column_id;
680 let inserted = seen.insert(col_id);
681 debug_assert!(
684 inserted,
685 "projection contains duplicate column id: {}",
686 col_id
687 );
688 read_col_ids.push(col_id);
690 }
691
692 if projection.is_empty() {
693 let time_index = metadata.time_index_column().column_id;
694 if seen.insert(time_index) {
695 read_col_ids.push(time_index);
696 }
697 }
698
699 let mut extra_col_names = HashSet::new();
700 let mut cols = HashSet::new();
701
702 if let Some(p) = predicate.predicate_without_region() {
703 for expr in p.exprs() {
704 cols.clear();
705 if expr_to_columns(expr, &mut cols).is_err() {
706 continue;
707 }
708 extra_col_names.extend(cols.iter().map(|col| col.name.clone()));
709 }
710 }
711
712 if let Some(expr) = predicate.region_partition_expr() {
713 expr.collect_column_names(&mut extra_col_names);
714 }
715
716 if !extra_col_names.is_empty() {
717 for col in &metadata.column_metadatas {
718 if extra_col_names.remove(&col.column_schema.name) && !seen.contains(&col.column_id)
719 {
720 read_col_ids.push(col.column_id);
721 }
722 }
723 if !extra_col_names.is_empty() {
724 warn!(
725 "Some columns in filters are not found in region {}: {:?}",
726 metadata.region_id, extra_col_names
727 );
728 }
729 }
730 Ok(read_col_ids)
731 }
732
733 fn build_read_columns(&self, col_ids: &[ColumnId]) -> Result<ReadColumns> {
739 let metadata = &self.version.metadata;
740 let json_type_hint = &self.request.json_type_hint;
741
742 let has_json2 = metadata
743 .schema
744 .arrow_schema()
745 .fields()
746 .iter()
747 .any(is_json2_extension_type);
748
749 if !has_json2 && json_type_hint.is_empty() {
750 return Ok(ReadColumns::new(col_ids.iter().copied()));
751 }
752
753 let mut json_target_types = BTreeMap::new();
754 for &col_id in col_ids {
755 let Some(col) = metadata.column_by_id(col_id) else {
756 continue;
757 };
758 let col_name = &col.column_schema.name;
759 let hint = json_type_hint.get(col_name);
760 if !col.column_schema.data_type.is_json2() {
761 ensure!(
762 hint.is_none(),
763 InvalidRequestSnafu {
764 region_id: metadata.region_id,
765 reason: format!(
766 "JSON type hint targets non-JSON2 column {} (id: {}, type: {})",
767 col_name, col_id, col.column_schema.data_type
768 ),
769 }
770 );
771 continue;
772 }
773 let target_type = hint.cloned().unwrap_or(JsonNativeType::Variant);
774 json_target_types.insert(col_id, target_type);
775 }
776 Ok(ReadColumns::new(col_ids.iter().copied()).with_json_target_types(json_target_types))
777 }
778
779 fn region_id(&self) -> RegionId {
780 self.version.metadata.region_id
781 }
782
783 fn build_time_range_predicate(&self) -> TimestampRange {
785 let time_index = self.version.metadata.time_index_column();
786 let unit = time_index
787 .column_schema
788 .data_type
789 .as_timestamp()
790 .expect("Time index must have timestamp-compatible type")
791 .unit();
792 build_time_range_predicate(&time_index.column_schema.name, unit, &self.request.filters)
793 }
794
795 fn partition_by_field_filters(&self) -> (Vec<Expr>, Vec<Expr>) {
798 let field_columns = self
799 .version
800 .metadata
801 .field_columns()
802 .map(|col| &col.column_schema.name)
803 .collect::<HashSet<_>>();
804
805 let mut columns = HashSet::new();
806
807 self.request.filters.iter().cloned().partition(|expr| {
808 columns.clear();
809 if expr_to_columns(expr, &mut columns).is_err() {
811 return true;
813 }
814 !columns
816 .iter()
817 .any(|column| field_columns.contains(&column.name))
818 })
819 }
820
821 fn build_invereted_index_applier(&self, filters: &[Expr]) -> Option<InvertedIndexApplierRef> {
823 if self.ignore_inverted_index {
824 return None;
825 }
826
827 let file_cache = self.cache_strategy.write_cache().map(|w| w.file_cache());
828 let inverted_index_cache = self.cache_strategy.inverted_index_cache().cloned();
829
830 let puffin_metadata_cache = self.cache_strategy.puffin_metadata_cache().cloned();
831
832 InvertedIndexApplierBuilder::new(
833 self.access_layer.table_dir().to_string(),
834 self.access_layer.path_type(),
835 self.access_layer.object_store().clone(),
836 self.version.metadata.as_ref(),
837 self.version.metadata.inverted_indexed_column_ids(
838 self.version
839 .options
840 .index_options
841 .inverted_index
842 .ignore_column_ids
843 .iter(),
844 ),
845 self.access_layer.puffin_manager_factory().clone(),
846 )
847 .with_file_cache(file_cache)
848 .with_inverted_index_cache(inverted_index_cache)
849 .with_puffin_metadata_cache(puffin_metadata_cache)
850 .build(filters)
851 .inspect_err(|err| warn!(err; "Failed to build invereted index applier"))
852 .ok()
853 .flatten()
854 .map(Arc::new)
855 }
856
857 fn build_bloom_filter_applier(&self, filters: &[Expr]) -> Option<BloomFilterIndexApplierRef> {
859 if self.ignore_bloom_filter {
860 return None;
861 }
862
863 let file_cache = self.cache_strategy.write_cache().map(|w| w.file_cache());
864 let bloom_filter_index_cache = self.cache_strategy.bloom_filter_index_cache().cloned();
865 let puffin_metadata_cache = self.cache_strategy.puffin_metadata_cache().cloned();
866
867 BloomFilterIndexApplierBuilder::new(
868 self.access_layer.table_dir().to_string(),
869 self.access_layer.path_type(),
870 self.access_layer.object_store().clone(),
871 self.version.metadata.as_ref(),
872 self.access_layer.puffin_manager_factory().clone(),
873 )
874 .with_file_cache(file_cache)
875 .with_bloom_filter_index_cache(bloom_filter_index_cache)
876 .with_puffin_metadata_cache(puffin_metadata_cache)
877 .build(filters)
878 .inspect_err(|err| warn!(err; "Failed to build bloom filter index applier"))
879 .ok()
880 .flatten()
881 .map(Arc::new)
882 }
883
884 fn build_fulltext_index_applier(&self, filters: &[Expr]) -> Option<FulltextIndexApplierRef> {
886 if self.ignore_fulltext_index {
887 return None;
888 }
889
890 let file_cache = self.cache_strategy.write_cache().map(|w| w.file_cache());
891 let puffin_metadata_cache = self.cache_strategy.puffin_metadata_cache().cloned();
892 let bloom_filter_index_cache = self.cache_strategy.bloom_filter_index_cache().cloned();
893 FulltextIndexApplierBuilder::new(
894 self.access_layer.table_dir().to_string(),
895 self.access_layer.path_type(),
896 self.access_layer.object_store().clone(),
897 self.access_layer.puffin_manager_factory().clone(),
898 self.version.metadata.as_ref(),
899 )
900 .with_file_cache(file_cache)
901 .with_puffin_metadata_cache(puffin_metadata_cache)
902 .with_bloom_filter_cache(bloom_filter_index_cache)
903 .build(filters)
904 .inspect_err(|err| warn!(err; "Failed to build fulltext index applier"))
905 .ok()
906 .flatten()
907 .map(Arc::new)
908 }
909
910 #[cfg(feature = "vector_index")]
912 fn build_vector_index_applier(&self) -> Option<VectorIndexApplierRef> {
913 let vector_search = self.request.vector_search.as_ref()?;
914
915 let file_cache = self.cache_strategy.write_cache().map(|w| w.file_cache());
916 let puffin_metadata_cache = self.cache_strategy.puffin_metadata_cache().cloned();
917 let vector_index_cache = self.cache_strategy.vector_index_cache().cloned();
918
919 let applier = VectorIndexApplier::new(
920 self.access_layer.table_dir().to_string(),
921 self.access_layer.path_type(),
922 self.access_layer.object_store().clone(),
923 self.access_layer.puffin_manager_factory().clone(),
924 vector_search.column_id,
925 vector_search.query_vector.clone(),
926 vector_search.metric,
927 )
928 .with_file_cache(file_cache)
929 .with_puffin_metadata_cache(puffin_metadata_cache)
930 .with_vector_index_cache(vector_index_cache);
931
932 Some(Arc::new(applier))
933 }
934}
935
936fn file_in_range(file: &FileHandle, predicate: &TimestampRange) -> bool {
938 if predicate == &TimestampRange::min_to_max() {
939 return true;
940 }
941 let (start, end) = file.time_range();
943 let file_ts_range = TimestampRange::new_inclusive(Some(start), Some(end));
944 file_ts_range.intersects(predicate)
945}
946
947fn time_range_covers_file(time_range: Option<&TimestampRange>, file: &FileHandle) -> bool {
949 let Some(time_range) = time_range else {
950 return false;
951 };
952 let (start, end) = file.time_range();
953 time_range.contains(&start) && time_range.contains(&end)
954}
955
956pub struct ScanInput {
958 pub(crate) series_index: Option<SeriesIndexReadContext>,
960 ignore_range_index: bool,
962 access_layer: AccessLayerRef,
964 pub(crate) mapper: Arc<FlatProjectionMapper>,
966 pub(crate) read_cols: ReadColumns,
970 pub(crate) time_range: Option<TimestampRange>,
972 scan_analysis: Option<ScanAnalysis>,
974 pub(crate) predicate: PredicateGroup,
976 region_partition_expr: Option<PartitionExpr>,
978 pub(crate) memtables: Vec<MemRangeBuilder>,
980 pub(crate) files: Vec<FileHandle>,
982 primary_key_mapper: OnceLock<Arc<PrimaryKeyRangeMapper>>,
984 batch_size: usize,
986 pub(crate) cache_strategy: CacheStrategy,
988 ignore_file_not_found: bool,
990 pub(crate) max_concurrent_scan_files: usize,
992 pub(crate) scan_memory_pool: Arc<dyn MemoryPool>,
994 inverted_index_appliers: [Option<InvertedIndexApplierRef>; 2],
996 bloom_filter_index_appliers: [Option<BloomFilterIndexApplierRef>; 2],
997 fulltext_index_appliers: [Option<FulltextIndexApplierRef>; 2],
998 #[cfg(feature = "vector_index")]
1000 pub(crate) vector_index_applier: Option<VectorIndexApplierRef>,
1001 #[cfg(feature = "vector_index")]
1003 pub(crate) vector_index_k: Option<usize>,
1004 pub(crate) query_start: Option<Instant>,
1006 pub(crate) append_mode: bool,
1008 pub(crate) filter_deleted: bool,
1010 pub(crate) merge_mode: MergeMode,
1012 pub(crate) series_row_selector: Option<TimeSeriesRowSelector>,
1014 pub(crate) distribution: Option<TimeSeriesDistribution>,
1016 explain_flat_format: bool,
1018 pub(crate) snapshot_sequence: Option<SequenceNumber>,
1020 pub(crate) sequence_range: Option<SequenceRange>,
1025 pub(crate) compaction: bool,
1027 json2_rewrite_targets: Json2RewriteTargets,
1029 pub(crate) query_stat_counters: Option<RegionQueryStatCounters>,
1031 #[cfg(feature = "enterprise")]
1032 extension_ranges: Vec<BoxedExtensionRange>,
1033}
1034
1035struct ScanAnalysis {
1038 fingerprint: Option<ScanRequestFingerprint>,
1040 implied_time_range: Option<TimestampRange>,
1046}
1047
1048pub(crate) struct ScanInputBuilder {
1050 input: ScanInput,
1051}
1052
1053impl ScanInput {
1054 #[must_use]
1056 pub(crate) fn builder(
1057 access_layer: AccessLayerRef,
1058 mapper: FlatProjectionMapper,
1059 ) -> ScanInputBuilder {
1060 ScanInputBuilder {
1061 input: ScanInput {
1062 series_index: None,
1063 ignore_range_index: false,
1064 access_layer,
1065 read_cols: mapper.read_columns().clone(),
1066 mapper: Arc::new(mapper),
1067 time_range: None,
1068 scan_analysis: None,
1069 predicate: PredicateGroup::default(),
1070 region_partition_expr: None,
1071 memtables: Vec::new(),
1072 files: Vec::new(),
1073 primary_key_mapper: OnceLock::new(),
1074 batch_size: crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
1075 cache_strategy: CacheStrategy::Disabled,
1076 ignore_file_not_found: false,
1077 max_concurrent_scan_files: DEFAULT_MAX_CONCURRENT_SCAN_FILES,
1078 scan_memory_pool: Arc::new(UnboundedMemoryPool::default()),
1079 inverted_index_appliers: [None, None],
1080 bloom_filter_index_appliers: [None, None],
1081 fulltext_index_appliers: [None, None],
1082 #[cfg(feature = "vector_index")]
1083 vector_index_applier: None,
1084 #[cfg(feature = "vector_index")]
1085 vector_index_k: None,
1086 query_start: None,
1087 append_mode: false,
1088 filter_deleted: true,
1089 merge_mode: MergeMode::default(),
1090 series_row_selector: None,
1091 distribution: None,
1092 explain_flat_format: false,
1093 snapshot_sequence: None,
1094 sequence_range: None,
1095 compaction: false,
1096 json2_rewrite_targets: Arc::default(),
1097 query_stat_counters: None,
1098 #[cfg(feature = "enterprise")]
1099 extension_ranges: Vec::new(),
1100 },
1101 }
1102 }
1103
1104 pub(crate) fn batch_size(&self) -> usize {
1106 self.batch_size
1107 }
1108
1109 pub(crate) fn primary_key_mapper(&self) -> &PrimaryKeyRangeMapper {
1111 self.primary_key_mapper
1112 .get_or_init(|| Arc::new(PrimaryKeyRangeMapper::new(self.region_metadata().clone())))
1113 }
1114
1115 pub(crate) fn implied_time_range(&self) -> Option<&TimestampRange> {
1117 self.scan_analysis
1118 .as_ref()
1119 .expect("ScanInput must be built")
1120 .implied_time_range
1121 .as_ref()
1122 }
1123
1124 pub(crate) fn scan_fingerprint(&self) -> Option<&ScanRequestFingerprint> {
1126 self.scan_analysis
1127 .as_ref()
1128 .expect("ScanInput must be built")
1129 .fingerprint
1130 .as_ref()
1131 }
1132}
1133
1134impl ScanInputBuilder {
1135 fn with_primary_key_mapper(mut self, mapper: Arc<PrimaryKeyRangeMapper>) -> Self {
1136 self.input.primary_key_mapper = OnceLock::from(mapper);
1137 self
1138 }
1139
1140 #[must_use]
1142 pub(crate) fn with_ignore_range_index(mut self, ignore: bool) -> Self {
1143 self.input.ignore_range_index = ignore;
1144 self
1145 }
1146
1147 #[must_use]
1149 pub(crate) fn with_series_index(
1150 mut self,
1151 series_index: Option<SeriesIndexReadContext>,
1152 ) -> Self {
1153 self.input.series_index = series_index;
1154 self
1155 }
1156
1157 #[must_use]
1159 pub(crate) fn with_time_range(mut self, time_range: Option<TimestampRange>) -> Self {
1160 self.input.time_range = time_range;
1161 self
1162 }
1163
1164 #[must_use]
1166 pub(crate) fn with_predicate(mut self, predicate: PredicateGroup) -> Self {
1167 self.input.region_partition_expr = predicate.region_partition_expr().cloned();
1168 self.input.predicate = predicate;
1169 self
1170 }
1171
1172 #[must_use]
1174 pub(crate) fn build(mut self) -> ScanInput {
1175 let input = &self.input;
1176 let eligible = !input.compaction
1177 && !input.files.is_empty()
1178 && matches!(input.cache_strategy, CacheStrategy::EnableAll(_));
1179
1180 let metadata = input.region_metadata();
1181 let tag_names: HashSet<&str> = metadata
1182 .column_metadatas
1183 .iter()
1184 .filter(|col| col.semantic_type == SemanticType::Tag)
1185 .map(|col| col.column_schema.name.as_str())
1186 .collect();
1187
1188 let time_index = metadata.time_index_column();
1189 let time_index_name = time_index.column_schema.name.clone();
1190 let ts_col_unit = time_index
1191 .column_schema
1192 .data_type
1193 .as_timestamp()
1194 .expect("Time index must have timestamp-compatible type")
1195 .unit();
1196
1197 let exprs = input
1198 .predicate_group()
1199 .predicate_without_region()
1200 .map(|predicate| predicate.exprs())
1201 .unwrap_or_default();
1202
1203 let mut filters = Vec::new();
1204 let mut time_only_exprs: Vec<&Expr> = Vec::new();
1205 let mut has_tag_filter = false;
1206 let mut columns = HashSet::new();
1207
1208 for expr in exprs {
1209 columns.clear();
1210 let is_time_only = match expr_to_columns(expr, &mut columns) {
1211 Ok(()) if !columns.is_empty() => {
1212 has_tag_filter |= columns
1213 .iter()
1214 .any(|col| tag_names.contains(col.name.as_str()));
1215 columns.iter().all(|col| col.name == time_index_name)
1216 }
1217 _ => false,
1218 };
1219
1220 if is_time_only
1228 && extract_time_range_from_expr(&time_index_name, ts_col_unit, expr).is_some()
1229 {
1230 time_only_exprs.push(expr);
1231 } else {
1232 filters.push(expr.to_string());
1233 }
1234 }
1235
1236 let implied_time_range =
1237 implied_time_range_from_exprs(&time_index_name, ts_col_unit, &time_only_exprs);
1238
1239 let fingerprint = if eligible && has_tag_filter {
1241 let mut time_filters: Vec<String> =
1242 time_only_exprs.iter().map(|e| e.to_string()).collect();
1243
1244 filters.sort_unstable();
1246 time_filters.sort_unstable();
1247 let read_columns = input.read_cols.clone();
1248 let fingerprint = crate::read::range_cache::ScanRequestFingerprintBuilder {
1249 read_column_types: read_columns
1250 .column_ids_iter()
1251 .map(|id| {
1252 read_columns
1253 .json_target_type(id)
1254 .cloned()
1255 .map(ConcreteDataType::json2)
1256 .or_else(|| {
1257 metadata
1258 .column_by_id(id)
1259 .map(|col| col.column_schema.data_type.clone())
1260 })
1261 })
1262 .collect(),
1263 read_columns,
1264 filters,
1265 time_filters,
1266 series_row_selector: input.series_row_selector,
1267 append_mode: input.append_mode,
1268 filter_deleted: input.filter_deleted,
1269 merge_mode: input.merge_mode,
1270 sequence_range: input.sequence_range,
1271 partition_expr_version: metadata.partition_expr_version,
1272 }
1273 .build();
1274 Some(fingerprint)
1275 } else {
1276 None
1277 };
1278
1279 self.input.scan_analysis = Some(ScanAnalysis {
1280 fingerprint,
1281 implied_time_range,
1282 });
1283 self.input
1284 }
1285
1286 #[must_use]
1288 pub(crate) fn with_memtables(mut self, memtables: Vec<MemRangeBuilder>) -> Self {
1289 self.input.memtables = memtables;
1290 self
1291 }
1292
1293 #[must_use]
1295 pub(crate) fn with_files(mut self, files: Vec<FileHandle>) -> Self {
1296 self.input.files = files;
1297 self
1298 }
1299
1300 #[must_use]
1302 pub(crate) fn with_batch_size(mut self, batch_size: usize) -> Self {
1303 self.input.batch_size = batch_size;
1304 self
1305 }
1306
1307 #[must_use]
1309 pub(crate) fn with_cache(mut self, cache: CacheStrategy) -> Self {
1310 self.input.cache_strategy = cache;
1311 self
1312 }
1313
1314 #[must_use]
1316 pub(crate) fn with_ignore_file_not_found(mut self, ignore: bool) -> Self {
1317 self.input.ignore_file_not_found = ignore;
1318 self
1319 }
1320
1321 #[must_use]
1323 pub(crate) fn with_max_concurrent_scan_files(
1324 mut self,
1325 max_concurrent_scan_files: usize,
1326 ) -> Self {
1327 self.input.max_concurrent_scan_files = max_concurrent_scan_files;
1328 self
1329 }
1330
1331 #[must_use]
1333 pub(crate) fn with_scan_memory_pool(mut self, scan_memory_pool: Arc<dyn MemoryPool>) -> Self {
1334 self.input.scan_memory_pool = scan_memory_pool;
1335 self
1336 }
1337
1338 #[must_use]
1340 pub(crate) fn with_inverted_index_appliers(
1341 mut self,
1342 appliers: [Option<InvertedIndexApplierRef>; 2],
1343 ) -> Self {
1344 self.input.inverted_index_appliers = appliers;
1345 self
1346 }
1347
1348 #[must_use]
1350 pub(crate) fn with_bloom_filter_index_appliers(
1351 mut self,
1352 appliers: [Option<BloomFilterIndexApplierRef>; 2],
1353 ) -> Self {
1354 self.input.bloom_filter_index_appliers = appliers;
1355 self
1356 }
1357
1358 #[must_use]
1360 pub(crate) fn with_fulltext_index_appliers(
1361 mut self,
1362 appliers: [Option<FulltextIndexApplierRef>; 2],
1363 ) -> Self {
1364 self.input.fulltext_index_appliers = appliers;
1365 self
1366 }
1367
1368 #[cfg(feature = "vector_index")]
1370 #[must_use]
1371 pub(crate) fn with_vector_index_applier(
1372 mut self,
1373 applier: Option<VectorIndexApplierRef>,
1374 ) -> Self {
1375 self.input.vector_index_applier = applier;
1376 self
1377 }
1378
1379 #[cfg(feature = "vector_index")]
1381 #[must_use]
1382 pub(crate) fn with_vector_index_k(mut self, k: Option<usize>) -> Self {
1383 self.input.vector_index_k = k;
1384 self
1385 }
1386
1387 #[must_use]
1389 pub(crate) fn with_start_time(mut self, now: Option<Instant>) -> Self {
1390 self.input.query_start = now;
1391 self
1392 }
1393
1394 #[must_use]
1395 pub(crate) fn with_append_mode(mut self, is_append_mode: bool) -> Self {
1396 self.input.append_mode = is_append_mode;
1397 self
1398 }
1399
1400 pub(crate) fn with_query_stat_counters(
1401 mut self,
1402 counters: Option<RegionQueryStatCounters>,
1403 ) -> Self {
1404 self.input.query_stat_counters = counters;
1405 self
1406 }
1407
1408 #[must_use]
1410 pub(crate) fn with_filter_deleted(mut self, filter_deleted: bool) -> Self {
1411 self.input.filter_deleted = filter_deleted;
1412 self
1413 }
1414
1415 #[must_use]
1417 pub(crate) fn with_merge_mode(mut self, merge_mode: MergeMode) -> Self {
1418 self.input.merge_mode = merge_mode;
1419 self
1420 }
1421
1422 #[must_use]
1424 pub(crate) fn with_distribution(
1425 mut self,
1426 distribution: Option<TimeSeriesDistribution>,
1427 ) -> Self {
1428 self.input.distribution = distribution;
1429 self
1430 }
1431
1432 #[must_use]
1434 pub(crate) fn with_explain_flat_format(mut self, explain_flat_format: bool) -> Self {
1435 self.input.explain_flat_format = explain_flat_format;
1436 self
1437 }
1438
1439 #[must_use]
1441 pub(crate) fn with_series_row_selector(
1442 mut self,
1443 series_row_selector: Option<TimeSeriesRowSelector>,
1444 ) -> Self {
1445 self.input.series_row_selector = series_row_selector;
1446 self
1447 }
1448
1449 #[must_use]
1450 pub(crate) fn with_snapshot_sequence(
1451 mut self,
1452 snapshot_sequence: Option<SequenceNumber>,
1453 ) -> Self {
1454 self.input.snapshot_sequence = snapshot_sequence;
1455 self
1456 }
1457
1458 #[must_use]
1459 pub(crate) fn with_sequence_range(mut self, sequence_range: Option<SequenceRange>) -> Self {
1460 self.input.sequence_range = sequence_range;
1461 self
1462 }
1463
1464 #[must_use]
1466 pub(crate) fn with_compaction(mut self, compaction: bool) -> Self {
1467 self.input.compaction = compaction;
1468 self
1469 }
1470
1471 #[must_use]
1473 pub(crate) fn with_json2_rewrite_targets(mut self, targets: Json2RewriteTargets) -> Self {
1474 self.input.json2_rewrite_targets = targets;
1475 self
1476 }
1477
1478 #[cfg(feature = "enterprise")]
1479 #[must_use]
1480 pub(crate) fn with_extension_ranges(
1481 mut self,
1482 extension_ranges: Vec<BoxedExtensionRange>,
1483 ) -> Self {
1484 self.input.extension_ranges = extension_ranges;
1485 self
1486 }
1487}
1488
1489impl ScanInput {
1490 pub(crate) fn build_mem_ranges(&self, index: RowGroupIndex) -> SmallVec<[MemtableRange; 2]> {
1492 let memtable = &self.memtables[index.index];
1493 let mut ranges = SmallVec::new();
1494 memtable.build_ranges(index.row_group_index, &mut ranges);
1495 ranges
1496 }
1497
1498 pub(crate) fn predicate_for_file(&self, file: &FileHandle) -> Option<Predicate> {
1499 if self.should_skip_region_partition(file) {
1500 self.predicate.predicate_without_region().cloned()
1501 } else {
1502 self.predicate.predicate().cloned()
1503 }
1504 }
1505
1506 fn should_skip_region_partition(&self, file: &FileHandle) -> bool {
1507 match (
1508 self.region_partition_expr.as_ref(),
1509 file.meta_ref().partition_expr.as_ref(),
1510 ) {
1511 (Some(region_expr), Some(file_expr)) => region_expr == file_expr,
1512 _ => false,
1513 }
1514 }
1515
1516 fn try_file_level_pruning_stats(&self, file: &FileHandle) -> Option<FileLevelPruningStats> {
1521 let (ts_min, ts_max) = file.time_range();
1522 let time_index = self.mapper.metadata().time_index_column();
1523 let time_index_unit = time_index.column_schema.data_type.as_timestamp()?.unit();
1524
1525 let min_ts = ts_min.convert_to(time_index_unit)?;
1528 let max_ts = ts_max.convert_to_ceil(time_index_unit)?;
1529
1530 Some(FileLevelPruningStats {
1531 min_scalar: timestamp_to_scalar_value(time_index_unit, Some(min_ts.value())),
1532 max_scalar: timestamp_to_scalar_value(time_index_unit, Some(max_ts.value())),
1533 time_index_col_name: time_index.column_schema.name.clone(),
1534 })
1535 }
1536
1537 #[inline]
1542 pub(crate) fn can_manifest_prune_file(&self, file: &FileHandle) -> bool {
1543 let predicate = self.predicate_for_file(file);
1544 self.manifest_prunes_file(file, predicate.as_ref())
1545 }
1546
1547 fn manifest_prunes_file(&self, file: &FileHandle, predicate: Option<&Predicate>) -> bool {
1548 if let Some(pred) = predicate
1549 && !pred.is_empty()
1550 && let Some(file_level_stats) = self.try_file_level_pruning_stats(file)
1551 {
1552 let pruning_results = pred.prune_with_stats(
1553 &file_level_stats,
1554 self.mapper.metadata().schema.arrow_schema(),
1555 );
1556 pruning_results.first() == Some(&false)
1557 } else {
1558 false
1559 }
1560 }
1561
1562 #[tracing::instrument(
1567 skip_all,
1568 fields(
1569 region_id = %self.region_metadata().region_id,
1570 file_id = %file.file_id()
1571 )
1572 )]
1573 pub async fn prune_file(
1574 &self,
1575 file: &FileHandle,
1576 pre_filter_mode: PreFilterMode,
1577 reader_metrics: &mut ReaderMetrics,
1578 ) -> Result<FileRangeBuilder> {
1579 let predicate = self.predicate_for_file(file);
1580
1581 if self.manifest_prunes_file(file, predicate.as_ref()) {
1583 reader_metrics.filter_metrics.files_time_range_pruned += 1;
1584 return Ok(FileRangeBuilder::default());
1585 }
1586
1587 self.prune_file_after_manifest_check(file, pre_filter_mode, true, predicate, reader_metrics)
1588 .await
1589 }
1590
1591 pub(crate) async fn prune_file_after_manifest_check(
1599 &self,
1600 file: &FileHandle,
1601 pre_filter_mode: PreFilterMode,
1602 enable_predicate_prefilter: bool,
1603 predicate: Option<Predicate>,
1604 reader_metrics: &mut ReaderMetrics,
1605 ) -> Result<FileRangeBuilder> {
1606 let may_build_selective_row_selection = predicate.is_some();
1607 let postpone_time_index_filter = time_range_covers_file(self.implied_time_range(), file);
1608 let decode_pk_values = !self.compaction
1609 && self
1610 .mapper
1611 .read_columns()
1612 .column_ids_iter()
1613 .any(|column_id| self.mapper.metadata().primary_key.contains(&column_id));
1614 let reader = self
1615 .access_layer
1616 .read_sst(file.clone())
1617 .series_index(
1618 self.series_index
1619 .clone()
1620 .filter(|_| !self.ignore_range_index),
1621 )
1622 .predicate(predicate)
1623 .projection(Some(self.read_cols.clone()))
1624 .json2_rewrite_targets(self.json2_rewrite_targets.clone())
1625 .cache(self.cache_strategy.clone())
1626 .inverted_index_appliers(self.inverted_index_appliers.clone())
1627 .bloom_filter_index_appliers(self.bloom_filter_index_appliers.clone())
1628 .fulltext_index_appliers(self.fulltext_index_appliers.clone());
1629 let reader = reader.batch_size(self.batch_size);
1630 let reader = if !self.compaction && may_build_selective_row_selection {
1631 reader.deferred_optional_page_index()
1632 } else {
1633 reader
1634 };
1635 #[cfg(feature = "vector_index")]
1636 let reader = {
1637 let mut reader = reader;
1638 reader =
1639 reader.vector_index_applier(self.vector_index_applier.clone(), self.vector_index_k);
1640 reader
1641 };
1642 let res = reader
1643 .expected_metadata(Some(self.mapper.metadata().clone()))
1644 .compaction(self.compaction)
1645 .pre_filter_mode(pre_filter_mode)
1646 .enable_predicate_prefilter(enable_predicate_prefilter)
1647 .postpone_time_index_filter(postpone_time_index_filter)
1648 .decode_primary_key_values(decode_pk_values)
1649 .build_reader_input(reader_metrics)
1650 .await;
1651 let read_input = match res {
1652 Ok(x) => x,
1653 Err(e) => {
1654 if e.is_object_not_found() && self.ignore_file_not_found {
1655 error!(e; "File to scan does not exist, region_id: {}, file: {}", file.region_id(), file.file_id());
1656 return Ok(FileRangeBuilder::default());
1657 } else {
1658 return Err(e);
1659 }
1660 }
1661 };
1662
1663 let Some((mut file_range_ctx, selection)) = read_input else {
1664 return Ok(FileRangeBuilder::default());
1665 };
1666
1667 let need_compat = !compat::has_same_columns_and_pk_encoding(
1668 &self.mapper,
1669 file_range_ctx.read_format(),
1670 self.compaction,
1671 );
1672 if need_compat {
1673 let compat = FlatCompatBatch::try_new(
1676 &self.mapper,
1677 file_range_ctx.read_format(),
1678 self.compaction,
1679 )?;
1680 file_range_ctx.set_compat_batch(compat);
1681 }
1682 Ok(FileRangeBuilder::new(Arc::new(file_range_ctx), selection))
1683 }
1684
1685 #[tracing::instrument(
1689 skip(self, sources, semaphore),
1690 fields(
1691 region_id = %self.region_metadata().region_id,
1692 source_count = sources.len()
1693 )
1694 )]
1695 pub(crate) fn create_parallel_flat_sources(
1696 &self,
1697 sources: Vec<BoxedRecordBatchStream>,
1698 semaphore: Arc<Semaphore>,
1699 channel_size: usize,
1700 ) -> Result<Vec<BoxedRecordBatchStream>> {
1701 if sources.len() <= 1 {
1702 return Ok(sources);
1703 }
1704
1705 let sources = sources
1707 .into_iter()
1708 .map(|source| {
1709 let (sender, receiver) = mpsc::channel(channel_size);
1710 self.spawn_flat_scan_task(source, semaphore.clone(), sender);
1711 let stream = Box::pin(ReceiverStream::new(receiver));
1712 Box::pin(stream) as _
1713 })
1714 .collect();
1715 Ok(sources)
1716 }
1717
1718 #[tracing::instrument(
1720 skip(self, input, semaphore, sender),
1721 fields(region_id = %self.region_metadata().region_id)
1722 )]
1723 pub(crate) fn spawn_flat_scan_task(
1724 &self,
1725 mut input: BoxedRecordBatchStream,
1726 semaphore: Arc<Semaphore>,
1727 sender: mpsc::Sender<Result<RecordBatch>>,
1728 ) {
1729 let region_id = self.region_metadata().region_id;
1730 let span = tracing::info_span!(
1731 "ScanInput::parallel_scan_task",
1732 region_id = %region_id,
1733 stream_kind = "flat"
1734 );
1735 common_runtime::spawn_query(
1736 async move {
1737 loop {
1738 let maybe_batch = {
1741 let _permit = semaphore.acquire().await.unwrap();
1743 input.next().await
1744 };
1745 match maybe_batch {
1746 Some(Ok(batch)) => {
1747 let _ = sender.send(Ok(batch)).await;
1748 }
1749 Some(Err(e)) => {
1750 let _ = sender.send(Err(e)).await;
1751 break;
1752 }
1753 None => break,
1754 }
1755 }
1756 }
1757 .instrument(span),
1758 );
1759 }
1760
1761 pub(crate) fn total_rows_is_exact(&self) -> bool {
1764 if self.region_partition_expr.is_none() {
1765 return true;
1766 }
1767
1768 #[cfg(feature = "enterprise")]
1770 if !self.extension_ranges.is_empty() {
1771 return false;
1772 }
1773
1774 self.files
1778 .iter()
1779 .all(|file| self.should_skip_region_partition(file))
1780 }
1781
1782 pub(crate) fn total_rows(&self) -> usize {
1783 let rows_in_files: usize = self.files.iter().map(|f| f.num_rows()).sum();
1784 let rows_in_memtables: usize = self.memtables.iter().map(|m| m.stats().num_rows()).sum();
1785
1786 let rows = rows_in_files + rows_in_memtables;
1787 #[cfg(feature = "enterprise")]
1788 let rows = rows
1789 + self
1790 .extension_ranges
1791 .iter()
1792 .map(|x| x.num_rows())
1793 .sum::<u64>() as usize;
1794 rows
1795 }
1796
1797 pub(crate) fn predicate_group(&self) -> &PredicateGroup {
1798 &self.predicate
1799 }
1800
1801 pub(crate) fn num_memtables(&self) -> usize {
1803 self.memtables.len()
1804 }
1805
1806 pub(crate) fn num_files(&self) -> usize {
1808 self.files.len()
1809 }
1810
1811 pub(crate) fn file_from_index(&self, index: RowGroupIndex) -> &FileHandle {
1813 let file_index = index.index - self.num_memtables();
1814 &self.files[file_index]
1815 }
1816
1817 pub fn region_metadata(&self) -> &RegionMetadataRef {
1818 self.mapper.metadata()
1819 }
1820
1821 fn range_pre_filter_mode(&self, source_count: usize) -> PreFilterMode {
1822 if source_count <= 1 {
1823 return PreFilterMode::All;
1828 }
1829
1830 pre_filter_mode(self.append_mode, self.merge_mode)
1831 }
1832}
1833
1834#[cfg(feature = "enterprise")]
1835impl ScanInput {
1836 #[cfg(feature = "enterprise")]
1837 pub(crate) fn extension_ranges(&self) -> &[BoxedExtensionRange] {
1838 &self.extension_ranges
1839 }
1840
1841 #[cfg(feature = "enterprise")]
1843 pub(crate) fn extension_range(&self, i: usize) -> &BoxedExtensionRange {
1844 &self.extension_ranges[i - self.num_memtables() - self.num_files()]
1845 }
1846}
1847
1848pub(crate) struct FileLevelPruningStats {
1852 pub(crate) min_scalar: ScalarValue,
1854 pub(crate) max_scalar: ScalarValue,
1856 pub(crate) time_index_col_name: String,
1858}
1859
1860impl PruningStatistics for FileLevelPruningStats {
1861 fn min_values(&self, column: &Column) -> Option<ArrayRef> {
1862 if column.name == self.time_index_col_name {
1863 ScalarValue::iter_to_array(std::iter::once(self.min_scalar.clone())).ok()
1864 } else {
1865 None
1866 }
1867 }
1868
1869 fn max_values(&self, column: &Column) -> Option<ArrayRef> {
1870 if column.name == self.time_index_col_name {
1871 ScalarValue::iter_to_array(std::iter::once(self.max_scalar.clone())).ok()
1872 } else {
1873 None
1874 }
1875 }
1876
1877 fn num_containers(&self) -> usize {
1878 1
1879 }
1880
1881 fn null_counts(&self, column: &Column) -> Option<ArrayRef> {
1882 if column.name == self.time_index_col_name {
1883 Some(Arc::new(UInt64Array::from(vec![0u64])))
1885 } else {
1886 None
1887 }
1888 }
1889
1890 fn row_counts(&self) -> Option<ArrayRef> {
1891 None
1892 }
1893
1894 fn contained(&self, _column: &Column, _values: &HashSet<ScalarValue>) -> Option<BooleanArray> {
1895 None
1896 }
1897}
1898
1899#[cfg(test)]
1900impl ScanInput {
1901 pub(crate) fn file_ids(&self) -> Vec<crate::sst::file::RegionFileId> {
1903 self.files.iter().map(|file| file.file_id()).collect()
1904 }
1905
1906 pub(crate) fn index_ids(&self) -> Vec<crate::sst::file::RegionIndexId> {
1907 self.files.iter().map(|file| file.index_id()).collect()
1908 }
1909}
1910
1911fn pre_filter_mode(append_mode: bool, merge_mode: MergeMode) -> PreFilterMode {
1912 if append_mode {
1913 return PreFilterMode::All;
1914 }
1915
1916 match merge_mode {
1917 MergeMode::LastRow => PreFilterMode::SkipFields,
1918 MergeMode::LastNonNull => PreFilterMode::SkipFields,
1919 }
1920}
1921
1922pub(crate) fn exact_sequence_range(
1942 request: &ScanRequest,
1943 version: &crate::region::version::Version,
1944) -> Result<(Vec<FileHandle>, Option<SequenceRange>)> {
1945 if request.skip_sst_files {
1946 return Ok((Vec::new(), None));
1947 }
1948
1949 let time_index = version.metadata.time_index_column();
1950 let unit = time_index
1951 .column_schema
1952 .data_type
1953 .as_timestamp()
1954 .expect("Time index must have timestamp-compatible type")
1955 .unit();
1956 let time_range =
1957 build_time_range_predicate(&time_index.column_schema.name, unit, &request.filters);
1958 let min = request.memtable_min_sequence;
1959 let mut check_capability = request.exact_sequence_range && min.is_some();
1960 let sst_min_sequence = request.sst_min_sequence.and_then(NonZeroU64::new);
1961 let mut files = Vec::new();
1962 let mut files_allow_exact_range = true;
1963
1964 for file in version
1965 .ssts
1966 .levels()
1967 .iter()
1968 .flat_map(|level| level.files.values())
1969 .filter(|file| file_in_range(file, &time_range))
1970 {
1971 let meta = file.meta_ref();
1972 let selected = (!request.exact_sequence_range
1973 || min.is_none_or(|min| meta.sequence.is_none_or(|sequence| sequence.get() > min)))
1974 && match (sst_min_sequence, meta.sequence) {
1975 (Some(min_sequence), Some(file_sequence)) => file_sequence > min_sequence,
1976 (Some(_), None) | (None, _) => true,
1979 };
1980 if !selected {
1981 continue;
1982 }
1983
1984 if let Some(min) = min
1985 && check_capability
1986 {
1987 if meta.region_id != version.metadata.region_id && meta.sequence.is_none() {
1988 return RegionSequenceDomainBrokenSnafu {
1989 region_id: version.metadata.region_id,
1990 file_region_id: meta.region_id,
1991 file_id: meta.file_id,
1992 }
1993 .fail();
1994 }
1995 if meta.region_id == version.metadata.region_id
1996 && !file.is_effective_target_sequence_trusted(version.metadata.region_id)
1997 && meta.sequence.is_none_or(|barrier| barrier.get() > min)
1998 {
1999 files_allow_exact_range = false;
2003 check_capability = false;
2004 }
2005 }
2006 files.push(file.clone());
2007 }
2008
2009 let sequence_range = match (
2010 request.exact_sequence_range,
2011 min,
2012 version.options.preserve_row_sequence,
2013 request.memtable_max_sequence,
2014 files_allow_exact_range,
2015 ) {
2016 (true, Some(min), true, Some(max), true) => Some(SequenceRange::GtLtEq { min, max }),
2017 _ => None,
2018 };
2019 Ok((files, sequence_range))
2020}
2021
2022pub struct StreamContext {
2025 pub input: ScanInput,
2027 pub(crate) ranges: Vec<RangeMeta>,
2029 pub(crate) query_start: Instant,
2032}
2033
2034impl StreamContext {
2035 pub(crate) fn seq_scan_ctx(input: ScanInput) -> Self {
2037 let query_start = input.query_start.unwrap_or_else(Instant::now);
2038 let ranges = RangeMeta::seq_scan_ranges(&input);
2039 READ_SST_COUNT.observe(input.num_files() as f64);
2040 Self {
2041 input,
2042 ranges,
2043 query_start,
2044 }
2045 }
2046
2047 pub(crate) fn unordered_scan_ctx(input: ScanInput) -> Self {
2049 let query_start = input.query_start.unwrap_or_else(Instant::now);
2050 let ranges = RangeMeta::unordered_scan_ranges(&input);
2051 READ_SST_COUNT.observe(input.num_files() as f64);
2052 Self {
2053 input,
2054 ranges,
2055 query_start,
2056 }
2057 }
2058
2059 pub(crate) fn is_mem_range_index(&self, index: RowGroupIndex) -> bool {
2061 self.input.num_memtables() > index.index
2062 }
2063
2064 pub(crate) fn is_file_range_index(&self, index: RowGroupIndex) -> bool {
2065 !self.is_mem_range_index(index)
2066 && index.index < self.input.num_files() + self.input.num_memtables()
2067 }
2068
2069 pub(crate) fn range_pre_filter_mode(&self, part_range: &PartitionRange) -> PreFilterMode {
2070 let range_meta = &self.ranges[part_range.identifier];
2071 let source_count = range_meta.indices.len();
2072
2073 self.input.range_pre_filter_mode(source_count)
2074 }
2075
2076 pub(crate) fn partition_ranges(&self) -> Vec<PartitionRange> {
2078 self.ranges
2079 .iter()
2080 .enumerate()
2081 .map(|(idx, range_meta)| range_meta.new_partition_range(idx))
2082 .collect()
2083 }
2084
2085 pub(crate) fn format_for_explain(&self, verbose: bool, f: &mut fmt::Formatter) -> fmt::Result {
2087 let (mut num_mem_ranges, mut num_file_ranges, mut num_other_ranges) = (0, 0, 0);
2088 for range_meta in &self.ranges {
2089 for idx in &range_meta.row_group_indices {
2090 if self.is_mem_range_index(*idx) {
2091 num_mem_ranges += 1;
2092 } else if self.is_file_range_index(*idx) {
2093 num_file_ranges += 1;
2094 } else {
2095 num_other_ranges += 1;
2096 }
2097 }
2098 }
2099 if verbose {
2100 write!(f, "{{")?;
2101 }
2102 write!(
2103 f,
2104 r#""partition_count":{{"count":{}, "mem_ranges":{}, "files":{}, "file_ranges":{}"#,
2105 self.ranges.len(),
2106 num_mem_ranges,
2107 self.input.num_files(),
2108 num_file_ranges,
2109 )?;
2110 if num_other_ranges > 0 {
2111 write!(f, r#", "other_ranges":{}"#, num_other_ranges)?;
2112 }
2113 write!(f, "}}")?;
2114
2115 if let Some(selector) = &self.input.series_row_selector {
2116 write!(f, ", \"selector\":\"{}\"", selector)?;
2117 }
2118 if let Some(distribution) = &self.input.distribution {
2119 write!(f, ", \"distribution\":\"{}\"", distribution)?;
2120 }
2121
2122 if verbose {
2123 self.format_verbose_content(f)?;
2124 }
2125
2126 Ok(())
2127 }
2128
2129 fn format_verbose_content(&self, f: &mut fmt::Formatter) -> fmt::Result {
2130 struct FileWrapper<'a> {
2131 file: &'a FileHandle,
2132 }
2133
2134 impl fmt::Debug for FileWrapper<'_> {
2135 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
2136 let (start, end) = self.file.time_range();
2137 write!(
2138 f,
2139 r#"{{"file_id":"{}","time_range_start":"{}::{}","time_range_end":"{}::{}","rows":{},"size":{},"index_size":{}}}"#,
2140 self.file.file_id(),
2141 start.value(),
2142 start.unit(),
2143 end.value(),
2144 end.unit(),
2145 self.file.num_rows(),
2146 self.file.size(),
2147 self.file.index_size()
2148 )
2149 }
2150 }
2151
2152 struct InputWrapper<'a> {
2153 input: &'a ScanInput,
2154 }
2155
2156 #[cfg(feature = "enterprise")]
2157 impl InputWrapper<'_> {
2158 fn format_extension_ranges(&self, f: &mut fmt::Formatter) -> fmt::Result {
2159 if self.input.extension_ranges.is_empty() {
2160 return Ok(());
2161 }
2162
2163 let mut delimiter = "";
2164 write!(f, ", extension_ranges: [")?;
2165 for range in self.input.extension_ranges() {
2166 write!(f, "{}{:?}", delimiter, range)?;
2167 delimiter = ", ";
2168 }
2169 write!(f, "]")?;
2170 Ok(())
2171 }
2172 }
2173
2174 impl fmt::Debug for InputWrapper<'_> {
2175 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
2176 let output_schema = self.input.mapper.output_schema();
2177 if !output_schema.is_empty() {
2178 let names: Vec<_> = output_schema
2179 .column_schemas()
2180 .iter()
2181 .map(|col| &col.name)
2182 .collect();
2183 write!(f, ", \"projection\": {:?}", names)?;
2184 }
2185 if let Some(predicate) = &self.input.predicate.predicate() {
2186 if !predicate.exprs().is_empty() {
2187 let exprs: Vec<_> =
2188 predicate.exprs().iter().map(|e| e.to_string()).collect();
2189 write!(f, ", \"filters\": {:?}", exprs)?;
2190 }
2191 if !predicate.dyn_filters().is_empty() {
2192 let dyn_filters: Vec<_> = predicate
2193 .dyn_filters()
2194 .iter()
2195 .map(|f| format!("{}", f))
2196 .collect();
2197 write!(f, ", \"dyn_filters\": {:?}", dyn_filters)?;
2198 }
2199 }
2200 #[cfg(feature = "vector_index")]
2201 if let Some(vector_index_k) = self.input.vector_index_k {
2202 write!(f, ", \"vector_index_k\": {}", vector_index_k)?;
2203 }
2204 if !self.input.files.is_empty() {
2205 write!(f, ", \"files\": ")?;
2206 f.debug_list()
2207 .entries(self.input.files.iter().map(|file| FileWrapper { file }))
2208 .finish()?;
2209 }
2210 write!(f, ", \"flat_format\": {}", self.input.explain_flat_format)?;
2211 #[cfg(feature = "enterprise")]
2212 self.format_extension_ranges(f)?;
2213
2214 Ok(())
2215 }
2216 }
2217
2218 write!(f, "{:?}", InputWrapper { input: &self.input })
2219 }
2220
2221 pub(crate) fn add_dyn_filter_to_predicate(
2224 self: &Arc<Self>,
2225 filter_exprs: Vec<Arc<dyn datafusion::physical_plan::PhysicalExpr>>,
2226 ) -> Vec<bool> {
2227 let mut supported = Vec::with_capacity(filter_exprs.len());
2228 let filter_expr = filter_exprs
2229 .into_iter()
2230 .filter_map(|expr| {
2231 if let Ok(dyn_filter) = (expr as Arc<dyn std::any::Any + Send + Sync + 'static>)
2232 .downcast::<datafusion::physical_plan::expressions::DynamicFilterPhysicalExpr>()
2233 {
2234 supported.push(true);
2235 Some(dyn_filter)
2236 } else {
2237 supported.push(false);
2238 None
2239 }
2240 })
2241 .collect();
2242 self.input.predicate.add_dyn_filters(filter_expr);
2243 supported
2244 }
2245}
2246
2247#[derive(Clone, Default)]
2250pub struct PredicateGroup {
2251 time_filters: Option<Arc<Vec<SimpleFilterEvaluator>>>,
2252 predicate_all: Predicate,
2254 predicate_without_region: Predicate,
2256 region_partition_expr: Option<PartitionExpr>,
2258}
2259
2260impl PredicateGroup {
2261 pub fn new(metadata: &RegionMetadata, exprs: &[Expr]) -> Result<Self> {
2263 let mut combined_exprs = exprs.to_vec();
2264 let mut region_partition_expr = None;
2265
2266 if let Some(expr_json) = metadata.partition_expr.as_ref()
2267 && !expr_json.is_empty()
2268 && let Some(expr) = PartitionExpr::from_json_str(expr_json)
2269 .context(InvalidPartitionExprSnafu { expr: expr_json })?
2270 {
2271 let logical_expr = expr
2272 .try_as_logical_expr()
2273 .context(InvalidPartitionExprSnafu {
2274 expr: expr_json.clone(),
2275 })?;
2276
2277 combined_exprs.push(logical_expr);
2278 region_partition_expr = Some(expr);
2279 }
2280
2281 let mut time_filters = Vec::with_capacity(combined_exprs.len());
2282 let mut columns = HashSet::new();
2284 for expr in &combined_exprs {
2285 columns.clear();
2286 let Some(filter) = Self::expr_to_filter(expr, metadata, &mut columns) else {
2287 continue;
2288 };
2289 time_filters.push(filter);
2290 }
2291 let time_filters = if time_filters.is_empty() {
2292 None
2293 } else {
2294 Some(Arc::new(time_filters))
2295 };
2296
2297 let predicate_all = Predicate::new(combined_exprs);
2298 let predicate_without_region = Predicate::new(exprs.to_vec());
2299
2300 Ok(Self {
2301 time_filters,
2302 predicate_all,
2303 predicate_without_region,
2304 region_partition_expr,
2305 })
2306 }
2307
2308 pub(crate) fn time_filters(&self) -> Option<Arc<Vec<SimpleFilterEvaluator>>> {
2310 self.time_filters.clone()
2311 }
2312
2313 pub(crate) fn predicate(&self) -> Option<&Predicate> {
2315 if self.predicate_all.is_empty() {
2316 None
2317 } else {
2318 Some(&self.predicate_all)
2319 }
2320 }
2321
2322 pub(crate) fn predicate_without_region(&self) -> Option<&Predicate> {
2324 if self.predicate_without_region.is_empty() {
2325 None
2326 } else {
2327 Some(&self.predicate_without_region)
2328 }
2329 }
2330
2331 pub(crate) fn add_dyn_filters(&self, dyn_filters: Vec<Arc<DynamicFilterPhysicalExpr>>) {
2333 self.predicate_all.add_dyn_filters(dyn_filters.clone());
2334 self.predicate_without_region.add_dyn_filters(dyn_filters);
2335 }
2336
2337 pub(crate) fn clear_dyn_filters(&self) {
2339 self.predicate_all.clear_dyn_filters();
2340 self.predicate_without_region.clear_dyn_filters();
2341 }
2342
2343 pub(crate) fn region_partition_expr(&self) -> Option<&PartitionExpr> {
2345 self.region_partition_expr.as_ref()
2346 }
2347
2348 fn expr_to_filter(
2349 expr: &Expr,
2350 metadata: &RegionMetadata,
2351 columns: &mut HashSet<Column>,
2352 ) -> Option<SimpleFilterEvaluator> {
2353 columns.clear();
2354 expr_to_columns(expr, columns).ok()?;
2357 if columns.len() > 1 {
2358 return None;
2360 }
2361 let column = columns.iter().next()?;
2362 let column_meta = metadata.column_by_name(&column.name)?;
2363 if column_meta.semantic_type == SemanticType::Timestamp {
2364 SimpleFilterEvaluator::try_new(expr)
2365 } else {
2366 None
2367 }
2368 }
2369}
2370
2371#[cfg(test)]
2372mod tests {
2373 use std::collections::BTreeMap;
2374 use std::sync::Arc;
2375
2376 use common_time::timestamp::{TimeUnit, Timestamp};
2377 use datafusion::physical_plan::expressions::{
2378 binary as physical_binary, col as physical_col, lit as physical_lit,
2379 };
2380 use datafusion_common::ScalarValue;
2381 use datafusion_expr::{Operator, col, lit};
2382 use datatypes::arrow::datatypes::{
2383 DataType as ArrowDataType, Field, Schema as ArrowSchema, TimeUnit as ArrowTimeUnit,
2384 };
2385 use datatypes::prelude::ConcreteDataType;
2386 use datatypes::schema::ColumnSchema;
2387 use datatypes::types::json_type::JsonObjectType;
2388 use datatypes::value::Value;
2389 use partition::expr::col as partition_col;
2390 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
2391 use store_api::storage::{RegionId, TimeSeriesDistribution, TimeSeriesRowSelector};
2392
2393 use super::*;
2394 use crate::cache::CacheManager;
2395 use crate::read::range_cache::ScanRequestFingerprintBuilder;
2396 use crate::sst::file::FileMeta;
2397 use crate::test_util::memtable_util::metadata_with_primary_key;
2398 use crate::test_util::scheduler_util::SchedulerEnv;
2399
2400 async fn new_scan_input(metadata: RegionMetadataRef, filters: Vec<Expr>) -> ScanInputBuilder {
2401 let env = SchedulerEnv::new().await;
2402 let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
2403 let predicate = PredicateGroup::new(metadata.as_ref(), &filters).unwrap();
2404 let file = FileHandle::new(
2405 crate::sst::file::FileMeta::default(),
2406 Arc::new(crate::sst::file_purger::NoopFilePurger),
2407 );
2408
2409 ScanInput::builder(env.access_layer.clone(), mapper)
2410 .with_predicate(predicate)
2411 .with_cache(CacheStrategy::EnableAll(Arc::new(
2412 CacheManager::builder()
2413 .range_result_cache_size(1024)
2414 .build(),
2415 )))
2416 .with_files(vec![file])
2417 }
2418
2419 #[tokio::test]
2420 async fn test_total_rows_is_exact_after_partition_filter() {
2421 let expr = partition_col("k0").gt_eq(Value::String("foo".into()));
2422 let other = partition_col("k0").gt_eq(Value::String("bar".into()));
2423 for (region_expr, file_exprs, exact) in [
2424 (None, vec![None], true),
2425 (None, vec![Some(expr.clone())], true),
2426 (Some(expr.clone()), vec![], true),
2427 (Some(expr.clone()), vec![Some(expr.clone())], true),
2428 (Some(expr.clone()), vec![None], false),
2429 (Some(expr.clone()), vec![Some(other)], false),
2430 (Some(expr.clone()), vec![Some(expr), None], false),
2431 ] {
2432 let mut builder =
2433 RegionMetadataBuilder::from_existing(metadata_with_primary_key(vec![0, 1], false));
2434 builder.partition_expr_json(region_expr.map(|expr| expr.as_json_str().unwrap()));
2435 let metadata = Arc::new(builder.build_without_validation().unwrap());
2436 let files = file_exprs
2437 .into_iter()
2438 .map(|partition_expr| {
2439 FileHandle::new(
2440 FileMeta {
2441 partition_expr,
2442 ..Default::default()
2443 },
2444 Arc::new(crate::sst::file_purger::NoopFilePurger),
2445 )
2446 })
2447 .collect();
2448 let input = new_scan_input(metadata, vec![])
2449 .await
2450 .with_files(files)
2451 .build();
2452 assert_eq!(input.total_rows_is_exact(), exact);
2453 }
2454 }
2455
2456 fn ts_lit(val: i64) -> datafusion_expr::Expr {
2458 lit(ScalarValue::TimestampMillisecond(Some(val), None))
2459 }
2460
2461 fn metadata_with_time_index_unit(unit: TimeUnit) -> RegionMetadataRef {
2462 let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 456));
2463 builder
2464 .push_column_metadata(ColumnMetadata {
2465 column_schema: ColumnSchema::new(
2466 "k0".to_string(),
2467 ConcreteDataType::string_datatype(),
2468 false,
2469 ),
2470 semantic_type: SemanticType::Tag,
2471 column_id: 0,
2472 })
2473 .push_column_metadata(ColumnMetadata {
2474 column_schema: ColumnSchema::new(
2475 "k1".to_string(),
2476 ConcreteDataType::uint32_datatype(),
2477 false,
2478 ),
2479 semantic_type: SemanticType::Tag,
2480 column_id: 1,
2481 })
2482 .push_column_metadata(ColumnMetadata {
2483 column_schema: ColumnSchema::new(
2484 "ts".to_string(),
2485 ConcreteDataType::timestamp_datatype(unit),
2486 false,
2487 ),
2488 semantic_type: SemanticType::Timestamp,
2489 column_id: 2,
2490 })
2491 .push_column_metadata(ColumnMetadata {
2492 column_schema: ColumnSchema::new(
2493 "v0".to_string(),
2494 ConcreteDataType::int64_datatype(),
2495 true,
2496 ),
2497 semantic_type: SemanticType::Field,
2498 column_id: 3,
2499 })
2500 .primary_key(vec![0, 1]);
2501
2502 Arc::new(builder.build().unwrap())
2503 }
2504
2505 fn file_handle_with_time_range(start: Timestamp, end: Timestamp) -> FileHandle {
2506 FileHandle::new(
2507 FileMeta {
2508 time_range: (start, end),
2509 ..Default::default()
2510 },
2511 Arc::new(crate::sst::file_purger::NoopFilePurger),
2512 )
2513 }
2514
2515 #[test]
2516 fn test_time_range_covers_file() {
2517 let file = file_handle_with_time_range(
2518 Timestamp::new_millisecond(1000),
2519 Timestamp::new_millisecond(2000),
2520 );
2521
2522 assert!(!time_range_covers_file(None, &file));
2523 assert!(time_range_covers_file(
2524 Some(&TimestampRange::min_to_max()),
2525 &file
2526 ));
2527 assert!(time_range_covers_file(
2528 Some(&TimestampRange::new_inclusive(
2529 Some(Timestamp::new_millisecond(1000)),
2530 Some(Timestamp::new_millisecond(2000)),
2531 )),
2532 &file
2533 ));
2534 assert!(time_range_covers_file(
2535 TimestampRange::with_unit(500, 3000, TimeUnit::Millisecond).as_ref(),
2536 &file
2537 ));
2538 assert!(!time_range_covers_file(
2539 TimestampRange::with_unit(1000, 2000, TimeUnit::Millisecond).as_ref(),
2540 &file
2541 ));
2542 assert!(!time_range_covers_file(
2543 TimestampRange::with_unit(1001, 3000, TimeUnit::Millisecond).as_ref(),
2544 &file
2545 ));
2546 assert!(!time_range_covers_file(
2547 Some(&TimestampRange::empty()),
2548 &file
2549 ));
2550
2551 let seconds_file = file_handle_with_time_range(
2552 Timestamp::new(1, TimeUnit::Second),
2553 Timestamp::new(2, TimeUnit::Second),
2554 );
2555 assert!(time_range_covers_file(
2556 TimestampRange::with_unit(1000, 2001, TimeUnit::Millisecond).as_ref(),
2557 &seconds_file
2558 ));
2559 }
2560
2561 #[tokio::test]
2562 async fn test_scan_input_builder_computes_scan_analysis() {
2563 let metadata = metadata_with_time_index_unit(TimeUnit::Millisecond);
2564 let env = SchedulerEnv::new().await;
2565 let file = file_handle_with_time_range(
2566 Timestamp::new_millisecond(1000),
2567 Timestamp::new_millisecond(2000),
2568 );
2569
2570 let make_input = |filters: &[Expr]| {
2571 let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
2572 ScanInput::builder(env.access_layer.clone(), mapper)
2573 .with_predicate(PredicateGroup::new(&metadata, filters).unwrap())
2574 .build()
2575 };
2576
2577 let covered = make_input(&[col("ts").gt_eq(ts_lit(500)), col("ts").lt(ts_lit(3000))]);
2578 assert!(time_range_covers_file(covered.implied_time_range(), &file));
2579
2580 let disjoint = make_input(&[col("ts").eq(ts_lit(1000)).or(col("ts").eq(ts_lit(2000)))]);
2583 assert!(disjoint.implied_time_range().is_none());
2584 assert!(!time_range_covers_file(
2585 disjoint.implied_time_range(),
2586 &file
2587 ));
2588
2589 let replaced = ScanInput::builder(
2590 env.access_layer.clone(),
2591 FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap(),
2592 )
2593 .with_predicate(PredicateGroup::new(&metadata, &[col("ts").gt(ts_lit(3000))]).unwrap())
2594 .build();
2595 assert!(!time_range_covers_file(
2596 replaced.implied_time_range(),
2597 &file
2598 ));
2599 }
2600
2601 #[tokio::test]
2602 async fn test_scan_input_uses_explicit_batch_size() {
2603 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
2604 let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
2605 let env = SchedulerEnv::new().await;
2606 let input = ScanInput::builder(env.access_layer.clone(), mapper).build();
2607 assert_eq!(
2608 crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
2609 input.batch_size()
2610 );
2611
2612 let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
2613 let input = ScanInput::builder(env.access_layer.clone(), mapper)
2614 .with_compaction(true)
2615 .with_batch_size(256)
2616 .build();
2617 assert_eq!(256, input.batch_size());
2618
2619 let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
2620 let input = ScanInput::builder(env.access_layer.clone(), mapper)
2621 .with_batch_size(256)
2622 .with_compaction(true)
2623 .with_compaction(false)
2624 .build();
2625 assert_eq!(256, input.batch_size());
2626 }
2627
2628 #[tokio::test]
2629 async fn test_build_scan_fingerprint_for_eligible_scan() {
2630 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
2631 let input = new_scan_input(
2632 metadata.clone(),
2633 vec![
2634 col("ts").gt_eq(ts_lit(1000)),
2635 col("k0").eq(lit("foo")),
2636 col("v0").gt(lit(1)),
2637 ],
2638 )
2639 .await
2640 .with_distribution(Some(TimeSeriesDistribution::PerSeries))
2641 .with_series_row_selector(Some(TimeSeriesRowSelector::LastRow { after_merge: true }))
2642 .with_merge_mode(MergeMode::LastNonNull)
2643 .with_filter_deleted(false)
2644 .build();
2645
2646 let fingerprint = input.scan_fingerprint().unwrap();
2647
2648 let expected = ScanRequestFingerprintBuilder {
2649 read_columns: input.read_cols.clone(),
2650 read_column_types: vec![
2651 metadata
2652 .column_by_id(0)
2653 .map(|col| col.column_schema.data_type.clone()),
2654 metadata
2655 .column_by_id(2)
2656 .map(|col| col.column_schema.data_type.clone()),
2657 metadata
2658 .column_by_id(3)
2659 .map(|col| col.column_schema.data_type.clone()),
2660 ],
2661 filters: vec![
2662 col("k0").eq(lit("foo")).to_string(),
2663 col("v0").gt(lit(1)).to_string(),
2664 ],
2665 time_filters: vec![col("ts").gt_eq(ts_lit(1000)).to_string()],
2666 series_row_selector: Some(TimeSeriesRowSelector::LastRow { after_merge: true }),
2667 append_mode: false,
2668 filter_deleted: false,
2669 merge_mode: MergeMode::LastNonNull,
2670 sequence_range: None,
2671 partition_expr_version: 0,
2672 }
2673 .build();
2674 assert_eq!(&expected, fingerprint);
2675 assert_eq!(
2676 input.series_row_selector,
2677 Some(TimeSeriesRowSelector::LastRow { after_merge: true })
2678 );
2679 }
2680
2681 #[tokio::test]
2682 async fn test_build_scan_fingerprint_requires_tag_filter() {
2683 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
2684 let input = new_scan_input(
2685 metadata,
2686 vec![col("ts").gt_eq(lit(1000)), col("v0").gt(lit(1))],
2687 )
2688 .await
2689 .build();
2690
2691 assert!(input.scan_fingerprint().is_none());
2692 }
2693
2694 #[tokio::test]
2695 async fn test_build_scan_fingerprint_respects_scan_eligibility() {
2696 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
2697 let filters = vec![col("k0").eq(lit("foo"))];
2698
2699 let disabled = ScanInput::builder(
2700 SchedulerEnv::new().await.access_layer.clone(),
2701 FlatProjectionMapper::new(&metadata, [0, 2, 3].into_iter()).unwrap(),
2702 )
2703 .with_predicate(PredicateGroup::new(metadata.as_ref(), &filters).unwrap())
2704 .build();
2705 assert!(disabled.scan_fingerprint().is_none());
2706
2707 let compaction = new_scan_input(metadata.clone(), filters.clone())
2708 .await
2709 .with_compaction(true)
2710 .build();
2711 assert!(compaction.scan_fingerprint().is_none());
2712
2713 let no_files = new_scan_input(metadata, filters)
2715 .await
2716 .with_files(vec![])
2717 .build();
2718 assert!(no_files.scan_fingerprint().is_none());
2719 }
2720
2721 #[tokio::test]
2722 async fn test_build_scan_fingerprint_tracks_schema_and_partition_expr_changes() {
2723 let base = metadata_with_primary_key(vec![0, 1], false);
2724 let mut builder = RegionMetadataBuilder::from_existing(base);
2725 let partition_expr = partition_col("k0")
2726 .gt_eq(Value::String("foo".into()))
2727 .as_json_str()
2728 .unwrap();
2729 builder.partition_expr_json(Some(partition_expr));
2730 let metadata = Arc::new(builder.build_without_validation().unwrap());
2731
2732 let input = new_scan_input(metadata.clone(), vec![col("k0").eq(lit("foo"))])
2733 .await
2734 .build();
2735 let fingerprint = input.scan_fingerprint().unwrap();
2736
2737 let expected = ScanRequestFingerprintBuilder {
2738 read_columns: input.read_cols.clone(),
2739 read_column_types: vec![
2740 metadata
2741 .column_by_id(0)
2742 .map(|col| col.column_schema.data_type.clone()),
2743 metadata
2744 .column_by_id(2)
2745 .map(|col| col.column_schema.data_type.clone()),
2746 metadata
2747 .column_by_id(3)
2748 .map(|col| col.column_schema.data_type.clone()),
2749 ],
2750 filters: vec![col("k0").eq(lit("foo")).to_string()],
2751 time_filters: vec![],
2752 series_row_selector: None,
2753 append_mode: false,
2754 filter_deleted: true,
2755 merge_mode: MergeMode::LastRow,
2756 sequence_range: None,
2757 partition_expr_version: metadata.partition_expr_version,
2758 }
2759 .build();
2760 assert_eq!(&expected, fingerprint);
2761 assert_ne!(0, metadata.partition_expr_version);
2762 }
2763
2764 #[tokio::test]
2765 async fn test_build_scan_fingerprint_uses_json_target_types() {
2766 let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 456));
2767 builder
2768 .push_column_metadata(ColumnMetadata {
2769 column_schema: ColumnSchema::new(
2770 "k0".to_string(),
2771 ConcreteDataType::string_datatype(),
2772 false,
2773 ),
2774 semantic_type: SemanticType::Tag,
2775 column_id: 0,
2776 })
2777 .push_column_metadata(ColumnMetadata {
2778 column_schema: ColumnSchema::new(
2779 "ts".to_string(),
2780 ConcreteDataType::timestamp_millisecond_datatype(),
2781 false,
2782 ),
2783 semantic_type: SemanticType::Timestamp,
2784 column_id: 1,
2785 })
2786 .push_column_metadata(ColumnMetadata {
2787 column_schema: ColumnSchema::new(
2788 "j".to_string(),
2789 ConcreteDataType::json2(JsonNativeType::Variant),
2790 true,
2791 ),
2792 semantic_type: SemanticType::Field,
2793 column_id: 2,
2794 })
2795 .primary_key(vec![0]);
2796 let metadata = Arc::new(builder.build().unwrap());
2797
2798 let make_input = |target_type| async {
2799 let env = SchedulerEnv::new().await;
2800 let read_cols = ReadColumns::new([0, 1, 2])
2801 .with_json_target_types(BTreeMap::from([(2, target_type)]));
2802 let mapper =
2803 FlatProjectionMapper::new_with_read_columns(&metadata, vec![0, 1, 2], read_cols)
2804 .unwrap();
2805 let predicate =
2806 PredicateGroup::new(metadata.as_ref(), &[col("k0").eq(lit("foo"))]).unwrap();
2807 let file = FileHandle::new(
2808 FileMeta::default(),
2809 Arc::new(crate::sst::file_purger::NoopFilePurger),
2810 );
2811 ScanInput::builder(env.access_layer.clone(), mapper)
2812 .with_predicate(predicate)
2813 .with_cache(CacheStrategy::EnableAll(Arc::new(
2814 CacheManager::builder()
2815 .range_result_cache_size(1024)
2816 .build(),
2817 )))
2818 .with_files(vec![file])
2819 .build()
2820 };
2821
2822 let int_target = JsonNativeType::i64();
2823 let string_target = JsonNativeType::String;
2824 let int_input = make_input(int_target.clone()).await;
2825 let string_input = make_input(string_target).await;
2826 let int_fingerprint = int_input.scan_fingerprint().unwrap();
2827 let string_fingerprint = string_input.scan_fingerprint().unwrap();
2828
2829 assert_ne!(int_fingerprint, string_fingerprint);
2830 assert_eq!(
2831 Some(&Some(ConcreteDataType::json2(int_target))),
2832 int_fingerprint.read_column_types().get(2)
2833 );
2834 }
2835
2836 #[tokio::test]
2837 async fn test_scan_input_rejects_json_type_hint_for_non_json2_column() {
2838 let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 456));
2839 builder
2840 .push_column_metadata(ColumnMetadata {
2841 column_schema: ColumnSchema::new(
2842 "k0".to_string(),
2843 ConcreteDataType::string_datatype(),
2844 false,
2845 ),
2846 semantic_type: SemanticType::Tag,
2847 column_id: 0,
2848 })
2849 .push_column_metadata(ColumnMetadata {
2850 column_schema: ColumnSchema::new(
2851 "ts".to_string(),
2852 ConcreteDataType::timestamp_millisecond_datatype(),
2853 false,
2854 ),
2855 semantic_type: SemanticType::Timestamp,
2856 column_id: 1,
2857 })
2858 .push_column_metadata(ColumnMetadata {
2859 column_schema: ColumnSchema::new(
2860 "j".to_string(),
2861 ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::from([(
2862 "a".to_string(),
2863 JsonNativeType::i64(),
2864 )]))),
2865 true,
2866 ),
2867 semantic_type: SemanticType::Field,
2868 column_id: 2,
2869 })
2870 .push_column_metadata(ColumnMetadata {
2871 column_schema: ColumnSchema::new(
2872 "v0".to_string(),
2873 ConcreteDataType::int64_datatype(),
2874 true,
2875 ),
2876 semantic_type: SemanticType::Field,
2877 column_id: 3,
2878 })
2879 .primary_key(vec![0]);
2880 let metadata = Arc::new(builder.build().unwrap());
2881 let mutable = Arc::new(crate::memtable::time_partition::TimePartitions::new(
2882 metadata.clone(),
2883 Arc::new(crate::test_util::memtable_util::EmptyMemtableBuilder::default()),
2884 0,
2885 None,
2886 ));
2887 let version = Arc::new(
2888 crate::region::version::VersionBuilder::new(metadata.clone(), mutable).build(),
2889 );
2890 let env = SchedulerEnv::new().await;
2891 let request = ScanRequest {
2892 projection: Some(vec![0, 1, 2, 3]),
2893 json_type_hint: std::collections::HashMap::from([(
2894 "v0".to_string(),
2895 JsonNativeType::i64(),
2896 )]),
2897 ..Default::default()
2898 };
2899
2900 let err = ScanRegion::new(
2901 version,
2902 env.access_layer.clone(),
2903 request,
2904 CacheStrategy::Disabled,
2905 )
2906 .scan_input()
2907 .await;
2908 let Err(err) = err else {
2909 panic!("scan input should reject JSON type hint for non-JSON2 column");
2910 };
2911
2912 assert!(err.to_string().contains("non-JSON2 column v0"));
2913 }
2914
2915 #[test]
2916 fn test_update_dyn_filters_with_empty_base_predicates() {
2917 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
2918 let predicate_group = PredicateGroup::new(metadata.as_ref(), &[]).unwrap();
2919 assert!(predicate_group.predicate().is_none());
2920 assert!(predicate_group.predicate_without_region().is_none());
2921
2922 let dyn_filter = Arc::new(DynamicFilterPhysicalExpr::new(vec![], physical_lit(false)));
2923 predicate_group.add_dyn_filters(vec![dyn_filter]);
2924
2925 let predicate_all = predicate_group.predicate().unwrap();
2926 assert!(predicate_all.exprs().is_empty());
2927 assert_eq!(1, predicate_all.dyn_filters().len());
2928
2929 let predicate_without_region = predicate_group.predicate_without_region().unwrap();
2930 assert!(predicate_without_region.exprs().is_empty());
2931 assert_eq!(1, predicate_without_region.dyn_filters().len());
2932 }
2933
2934 #[test]
2935 fn test_clear_dyn_filters_preserves_predicate_group_static_filters() {
2936 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
2937 let static_filters = vec![col("k0").eq(lit("foo"))];
2938 let predicate_group = PredicateGroup::new(metadata.as_ref(), &static_filters).unwrap();
2939 let dynamic_filter = Arc::new(DynamicFilterPhysicalExpr::new(vec![], physical_lit(true)));
2940 predicate_group.add_dyn_filters(vec![dynamic_filter.clone()]);
2941
2942 predicate_group.clear_dyn_filters();
2943 dynamic_filter.update(physical_lit(false)).unwrap();
2945
2946 for predicate in [
2947 predicate_group.predicate().unwrap(),
2948 predicate_group.predicate_without_region().unwrap(),
2949 ] {
2950 assert_eq!(predicate.exprs(), static_filters);
2951 assert!(predicate.dyn_filters().is_empty());
2952 }
2953 }
2954
2955 #[test]
2956 fn test_file_level_pruning_stats_prunes_old_file() {
2957 let ts_col_name = "ts";
2958 let predicate = Predicate::new(vec![col(ts_col_name).gt(ts_lit(1000))]);
2959 let arrow_schema = Arc::new(ArrowSchema::new(vec![Field::new(
2960 ts_col_name,
2961 ArrowDataType::Timestamp(ArrowTimeUnit::Millisecond, None),
2962 false,
2963 )]));
2964
2965 let stats = FileLevelPruningStats {
2967 min_scalar: ScalarValue::TimestampMillisecond(Some(0), None),
2968 max_scalar: ScalarValue::TimestampMillisecond(Some(500), None),
2969 time_index_col_name: ts_col_name.to_string(),
2970 };
2971 assert_eq!(
2972 vec![false],
2973 predicate.prune_with_stats(&stats, &arrow_schema)
2974 );
2975
2976 let stats = FileLevelPruningStats {
2978 min_scalar: ScalarValue::TimestampMillisecond(Some(0), None),
2979 max_scalar: ScalarValue::TimestampMillisecond(Some(2000), None),
2980 time_index_col_name: ts_col_name.to_string(),
2981 };
2982 assert_eq!(
2983 vec![true],
2984 predicate.prune_with_stats(&stats, &arrow_schema)
2985 );
2986 }
2987
2988 #[test]
2989 fn test_file_level_pruning_stats_no_predicate_keeps_all() {
2990 let predicate = Predicate::new(vec![]);
2991 assert!(predicate.is_empty());
2992
2993 let stats = FileLevelPruningStats {
2994 min_scalar: ScalarValue::TimestampMillisecond(Some(0), None),
2995 max_scalar: ScalarValue::TimestampMillisecond(Some(500), None),
2996 time_index_col_name: "ts".to_string(),
2997 };
2998 let arrow_schema = Arc::new(ArrowSchema::new(Vec::<Field>::new()));
2999 assert_eq!(
3000 vec![true],
3001 predicate.prune_with_stats(&stats, &arrow_schema)
3002 );
3003 }
3004
3005 #[tokio::test]
3006 async fn test_file_level_pruning_stats_ceil_max_unit_conversion() {
3007 let metadata = metadata_with_time_index_unit(TimeUnit::Millisecond);
3008 let input = new_scan_input(metadata, vec![]).await.build();
3009 let file = file_handle_with_time_range(
3010 Timestamp::new(1_000_001, TimeUnit::Nanosecond),
3011 Timestamp::new(1_000_001, TimeUnit::Nanosecond),
3012 );
3013
3014 let stats = input.try_file_level_pruning_stats(&file).unwrap();
3015 assert_eq!(
3016 ScalarValue::TimestampMillisecond(Some(1), None),
3017 stats.min_scalar
3018 );
3019 assert_eq!(
3020 ScalarValue::TimestampMillisecond(Some(2), None),
3021 stats.max_scalar
3022 );
3023
3024 let predicate = Predicate::new(vec![col("ts").gt(ts_lit(1))]);
3026 assert_eq!(
3027 vec![true],
3028 predicate.prune_with_stats(&stats, input.mapper.metadata().schema.arrow_schema())
3029 );
3030 }
3031
3032 #[tokio::test]
3033 async fn test_file_level_pruning_stats_overflow_keeps_file() {
3034 let metadata = metadata_with_time_index_unit(TimeUnit::Nanosecond);
3035 let input = new_scan_input(metadata, vec![]).await.build();
3036 let file = file_handle_with_time_range(
3037 Timestamp::new(0, TimeUnit::Second),
3038 Timestamp::new(i64::MAX, TimeUnit::Second),
3039 );
3040
3041 assert!(input.try_file_level_pruning_stats(&file).is_none());
3042 }
3043
3044 #[test]
3045 fn test_file_level_pruning_stats_keeps_inclusive_boundary() {
3046 let ts_col_name = "ts";
3047 let predicate = Predicate::new(vec![col(ts_col_name).gt_eq(ts_lit(1000))]);
3048 let arrow_schema = Arc::new(ArrowSchema::new(vec![Field::new(
3049 ts_col_name,
3050 ArrowDataType::Timestamp(ArrowTimeUnit::Millisecond, None),
3051 false,
3052 )]));
3053 let stats = FileLevelPruningStats {
3054 min_scalar: ScalarValue::TimestampMillisecond(Some(0), None),
3055 max_scalar: ScalarValue::TimestampMillisecond(Some(1000), None),
3056 time_index_col_name: ts_col_name.to_string(),
3057 };
3058
3059 assert_eq!(
3060 vec![true],
3061 predicate.prune_with_stats(&stats, &arrow_schema)
3062 );
3063 }
3064
3065 #[tokio::test]
3066 async fn test_file_level_pruning_with_dyn_filter_only_predicate() {
3067 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
3068 let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
3069 let predicate_group = PredicateGroup::new(metadata.as_ref(), &[]).unwrap();
3070 predicate_group.add_dyn_filters(vec![Arc::new(DynamicFilterPhysicalExpr::new(
3071 vec![],
3072 physical_lit(false),
3073 ))]);
3074 let input = ScanInput::builder(SchedulerEnv::new().await.access_layer.clone(), mapper)
3075 .with_predicate(predicate_group)
3076 .build();
3077 let file = file_handle_with_time_range(
3078 Timestamp::new_millisecond(0),
3079 Timestamp::new_millisecond(1000),
3080 );
3081 let mut reader_metrics = ReaderMetrics::default();
3082
3083 let builder = input
3084 .prune_file(&file, PreFilterMode::SkipFields, &mut reader_metrics)
3085 .await
3086 .unwrap();
3087
3088 assert_eq!(1, reader_metrics.filter_metrics.files_time_range_pruned);
3089 let mut ranges = SmallVec::new();
3090 builder.build_ranges(-1, &mut ranges);
3091 assert!(ranges.is_empty());
3092 }
3093
3094 #[tokio::test]
3095 async fn test_manifest_pruning_observes_dynamic_filter_update() {
3096 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
3097 let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
3098 let predicate_group = PredicateGroup::new(metadata.as_ref(), &[]).unwrap();
3099 let arrow_schema = metadata.schema.arrow_schema();
3100 let ts_expr = physical_col("ts", arrow_schema.as_ref()).unwrap();
3101 let dyn_filter = Arc::new(DynamicFilterPhysicalExpr::new(
3102 vec![ts_expr.clone()],
3103 physical_lit(true),
3104 ));
3105 predicate_group.add_dyn_filters(vec![dyn_filter.clone()]);
3106 let input = ScanInput::builder(SchedulerEnv::new().await.access_layer.clone(), mapper)
3107 .with_predicate(predicate_group)
3108 .build();
3109 let file = file_handle_with_time_range(
3110 Timestamp::new_millisecond(0),
3111 Timestamp::new_millisecond(1000),
3112 );
3113
3114 assert!(!input.can_manifest_prune_file(&file));
3115
3116 let updated = physical_binary(
3117 ts_expr,
3118 Operator::Gt,
3119 physical_lit(ScalarValue::TimestampMillisecond(Some(1000), None)),
3120 arrow_schema.as_ref(),
3121 )
3122 .unwrap();
3123 dyn_filter.update(updated).unwrap();
3124
3125 assert!(input.can_manifest_prune_file(&file));
3126 input.predicate.clear_dyn_filters();
3127 assert!(!input.can_manifest_prune_file(&file));
3128 }
3129
3130 #[tokio::test]
3131 async fn test_range_pre_filter_mode() {
3132 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
3133 let cases = [
3134 (true, MergeMode::LastRow, 1, PreFilterMode::All),
3135 (false, MergeMode::LastNonNull, 1, PreFilterMode::All),
3136 (false, MergeMode::LastRow, 2, PreFilterMode::SkipFields),
3137 (true, MergeMode::LastRow, 2, PreFilterMode::All),
3138 ];
3139
3140 for (append_mode, merge_mode, source_count, expected_mode) in cases {
3141 let input = new_scan_input(metadata.clone(), vec![])
3142 .await
3143 .with_append_mode(append_mode)
3144 .with_merge_mode(merge_mode)
3145 .build();
3146
3147 assert_eq!(expected_mode, input.range_pre_filter_mode(source_count));
3148 }
3149 }
3150
3151 #[test]
3152 fn test_exact_sequence_selection_preserves_range_and_fallback() {
3153 let request = ScanRequest {
3154 exact_sequence_range: true,
3155 memtable_min_sequence: Some(1),
3156 memtable_max_sequence: Some(2),
3157 ..Default::default()
3158 };
3159 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
3160 let mutable = Arc::new(crate::memtable::time_partition::TimePartitions::new(
3161 metadata.clone(),
3162 Arc::new(crate::test_util::memtable_util::EmptyMemtableBuilder::default()),
3163 0,
3164 None,
3165 ));
3166 let version = Arc::new(
3167 crate::region::version::VersionBuilder::new(metadata, mutable)
3168 .options(crate::region::options::RegionOptions {
3169 preserve_row_sequence: true,
3170 ..Default::default()
3171 })
3172 .build(),
3173 );
3174 let (files, range) = exact_sequence_range(&request, &version).unwrap();
3175 assert!(files.is_empty());
3176 assert_eq!(Some(SequenceRange::GtLtEq { min: 1, max: 2 }), range);
3177
3178 let skip_sst_request = ScanRequest {
3179 skip_sst_files: true,
3180 ..request
3181 };
3182 let (files, range) = exact_sequence_range(&skip_sst_request, &version).unwrap();
3183 assert!(files.is_empty());
3184 assert_eq!(None, range);
3185 }
3186
3187 #[test]
3188 fn test_exact_sequence_selection_checks_selected_files() {
3189 use std::num::NonZeroU64;
3190
3191 use crate::test_util::new_noop_file_purger;
3192
3193 let request = ScanRequest {
3194 exact_sequence_range: true,
3195 memtable_min_sequence: Some(5),
3196 memtable_max_sequence: Some(10),
3197 ..Default::default()
3198 };
3199 let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
3200 let mutable = Arc::new(crate::memtable::time_partition::TimePartitions::new(
3201 metadata.clone(),
3202 Arc::new(crate::test_util::memtable_util::EmptyMemtableBuilder::default()),
3203 0,
3204 None,
3205 ));
3206 let target_region_id = metadata.region_id;
3207
3208 for (sequence, marked, expected_files, expected_range) in [
3211 (NonZeroU64::new(5), false, false, true),
3212 (NonZeroU64::new(5), true, false, true),
3213 (NonZeroU64::new(100), false, true, false),
3214 (None, false, true, false),
3215 (None, true, true, true),
3216 ] {
3217 let version =
3218 crate::region::version::VersionBuilder::new(metadata.clone(), mutable.clone())
3219 .options(crate::region::options::RegionOptions {
3220 preserve_row_sequence: true,
3221 ..Default::default()
3222 })
3223 .add_files(
3224 new_noop_file_purger(),
3225 [FileMeta {
3226 region_id: target_region_id,
3227 sequence,
3228 preserve_row_sequence: marked,
3229 ..Default::default()
3230 }]
3231 .into_iter(),
3232 )
3233 .build();
3234 let (files, range) = exact_sequence_range(&request, &version).unwrap();
3235 assert_eq!(expected_files, !files.is_empty());
3236 assert_eq!(expected_range, range.is_some());
3237 }
3238
3239 for (sequence, marked, expected_count, expected_error) in [
3242 (NonZeroU64::new(5), false, 0, false),
3243 (NonZeroU64::new(6), false, 1, false),
3244 (NonZeroU64::new(11), true, 1, false),
3245 (None, true, 1, true),
3246 ] {
3247 let file_id = store_api::storage::FileId::random();
3248 let version =
3249 crate::region::version::VersionBuilder::new(metadata.clone(), mutable.clone())
3250 .options(crate::region::options::RegionOptions {
3251 preserve_row_sequence: true,
3252 ..Default::default()
3253 })
3254 .add_files(
3255 new_noop_file_purger(),
3256 [FileMeta {
3257 region_id: RegionId::new(1, 2),
3258 file_id,
3259 sequence,
3260 preserve_row_sequence: marked,
3261 ..Default::default()
3262 }]
3263 .into_iter(),
3264 )
3265 .build();
3266 let result = exact_sequence_range(&request, &version);
3267 if expected_error {
3268 assert!(result.is_err());
3269 } else {
3270 let (files, range) = result.unwrap();
3271 assert_eq!(expected_count, files.len());
3272 assert!(range.is_some());
3273 }
3274 }
3275
3276 let request = ScanRequest {
3279 memtable_max_sequence: None,
3280 ..request
3281 };
3282 let version = crate::region::version::VersionBuilder::new(metadata, mutable)
3283 .options(crate::region::options::RegionOptions {
3284 preserve_row_sequence: false,
3285 ..Default::default()
3286 })
3287 .add_files(
3288 new_noop_file_purger(),
3289 [FileMeta {
3290 region_id: RegionId::new(1, 2),
3291 ..Default::default()
3292 }]
3293 .into_iter(),
3294 )
3295 .build();
3296 assert!(exact_sequence_range(&request, &version).is_err());
3297 }
3298}