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