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