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