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