1use std::collections::HashSet;
22use std::ops::{BitAnd, Range};
23use std::sync::Arc;
24
25use api::v1::SemanticType;
26use common_recordbatch::filter::SimpleFilterEvaluator;
27use datafusion_common::ScalarValue;
28use datafusion_expr::Expr;
29use datatypes::arrow::array::{Array, BinaryArray, BooleanArray, BooleanBufferBuilder};
30use datatypes::arrow::buffer::BooleanBuffer;
31use datatypes::arrow::datatypes::SchemaRef;
32use datatypes::arrow::record_batch::RecordBatch;
33use datatypes::prelude::ConcreteDataType;
34use datatypes::value::Value;
35use futures::StreamExt;
36use mito_codec::row_converter::{PrimaryKeyCodec, PrimaryKeyFilter, build_primary_key_codec};
37use parquet::arrow::ProjectionMask;
38use parquet::arrow::arrow_reader::{RowSelection, RowSelector};
39use parquet::file::metadata::ParquetMetaData;
40use parquet::schema::types::SchemaDescriptor;
41use smallvec::{SmallVec, smallvec};
42use snafu::{OptionExt, ResultExt};
43use store_api::metadata::{RegionMetadata, RegionMetadataRef};
44use store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME;
45use table::predicate::Predicate;
46
47use crate::cache::PrefilterKey;
48use crate::error::{
49 ComputeArrowSnafu, DecodeSnafu, EvalPartitionFilterSnafu, NewRecordBatchSnafu,
50 RecordBatchSnafu, Result, UnexpectedSnafu,
51};
52use crate::sst::parquet::file_range::PreFilterMode;
53use crate::sst::parquet::flat_format::FlatReadFormat;
54use crate::sst::parquet::format::{PrimaryKeyArray, StatValues};
55use crate::sst::parquet::reader::{
56 MaybeFilter, PhysicalFilterContext, RowGroupBuildContext, RowGroupReaderBuilder,
57 SimpleFilterContext,
58};
59
60pub(crate) fn matching_row_ranges_by_primary_key(
61 input: &RecordBatch,
62 pk_column_index: usize,
63 pk_filter: &mut dyn PrimaryKeyFilter,
64) -> Result<Vec<Range<usize>>> {
65 let pk_column = input.column(pk_column_index);
66 if let Some(pk_dict_array) = pk_column.as_any().downcast_ref::<PrimaryKeyArray>() {
67 matching_row_ranges_from_dict(pk_dict_array, input.num_rows(), pk_filter)
68 } else if let Some(pk_binary_array) = pk_column.as_any().downcast_ref::<BinaryArray>() {
69 matching_row_ranges_from_binary(pk_binary_array, input.num_rows(), pk_filter)
70 } else {
71 UnexpectedSnafu {
72 reason: format!(
73 "Primary key column is neither a dictionary nor a binary array, got {:?}",
74 pk_column.data_type()
75 ),
76 }
77 .fail()
78 }
79}
80
81fn push_matched_range(
84 matched_row_ranges: &mut Vec<Range<usize>>,
85 pk_filter: &mut dyn PrimaryKeyFilter,
86 pk: &[u8],
87 start: usize,
88 end: usize,
89) -> Result<()> {
90 if pk_filter.matches(pk).context(DecodeSnafu)? {
91 if let Some(last) = matched_row_ranges.last_mut()
92 && last.end == start
93 {
94 last.end = end;
95 } else {
96 matched_row_ranges.push(start..end);
97 }
98 }
99 Ok(())
100}
101
102fn matching_row_ranges_from_dict(
104 pk_dict_array: &PrimaryKeyArray,
105 num_rows: usize,
106 pk_filter: &mut dyn PrimaryKeyFilter,
107) -> Result<Vec<Range<usize>>> {
108 let pk_values = pk_dict_array
109 .values()
110 .as_any()
111 .downcast_ref::<BinaryArray>()
112 .context(UnexpectedSnafu {
113 reason: "Primary key values are not binary array",
114 })?;
115 let keys = pk_dict_array.keys();
116 let key_values = keys.values();
117
118 if key_values.is_empty() {
119 return Ok(std::iter::once(0..num_rows).collect());
120 }
121
122 let mut matched_row_ranges: Vec<Range<usize>> = Vec::new();
123 let mut start = 0;
124 while start < key_values.len() {
125 let key = key_values[start];
126 let mut end = start + 1;
127 while end < key_values.len() && key_values[end] == key {
128 end += 1;
129 }
130
131 push_matched_range(
132 &mut matched_row_ranges,
133 pk_filter,
134 pk_values.value(key as usize),
135 start,
136 end,
137 )?;
138
139 start = end;
140 }
141
142 Ok(matched_row_ranges)
143}
144
145fn matching_row_ranges_from_binary(
151 pk_array: &BinaryArray,
152 num_rows: usize,
153 pk_filter: &mut dyn PrimaryKeyFilter,
154) -> Result<Vec<Range<usize>>> {
155 if pk_array.is_empty() {
156 return Ok(std::iter::once(0..num_rows).collect());
157 }
158
159 let mut matched_row_ranges: Vec<Range<usize>> = Vec::new();
160 let mut start = 0;
161 while start < pk_array.len() {
162 let value = pk_array.value(start);
163 let mut end = start + 1;
164 while end < pk_array.len() && pk_array.value(end) == value {
165 end += 1;
166 }
167
168 push_matched_range(&mut matched_row_ranges, pk_filter, value, start, end)?;
169
170 start = end;
171 }
172
173 Ok(matched_row_ranges)
174}
175
176pub(crate) fn prefilter_flat_batch_by_primary_key(
179 input: RecordBatch,
180 pk_column_index: usize,
181 pk_filter: &mut dyn PrimaryKeyFilter,
182) -> Result<Option<RecordBatch>> {
183 if input.num_rows() == 0 {
184 return Ok(Some(input));
185 }
186
187 let matched_row_ranges =
188 matching_row_ranges_by_primary_key(&input, pk_column_index, pk_filter)?;
189 if matched_row_ranges.is_empty() {
190 return Ok(None);
191 }
192
193 if matched_row_ranges.len() == 1
194 && matched_row_ranges[0].start == 0
195 && matched_row_ranges[0].end == input.num_rows()
196 {
197 return Ok(Some(input));
198 }
199
200 if matched_row_ranges.len() == 1 {
201 let span = &matched_row_ranges[0];
202 return Ok(Some(input.slice(span.start, span.end - span.start)));
203 }
204
205 let mut builder = BooleanBufferBuilder::new(input.num_rows());
206 builder.append_n(input.num_rows(), false);
207 for span in matched_row_ranges {
208 for i in span {
209 builder.set_bit(i, true);
210 }
211 }
212
213 let filtered = datatypes::arrow::compute::filter_record_batch(
214 &input,
215 &BooleanArray::new(builder.finish(), None),
216 )
217 .context(ComputeArrowSnafu)?;
218 if filtered.num_rows() == 0 {
219 Ok(None)
220 } else {
221 Ok(Some(filtered))
222 }
223}
224
225pub(crate) struct CachedPrimaryKeyFilter {
226 inner: Box<dyn PrimaryKeyFilter>,
227 last_primary_key: Vec<u8>,
228 last_match: Option<bool>,
229}
230
231impl CachedPrimaryKeyFilter {
232 pub(crate) fn new(inner: Box<dyn PrimaryKeyFilter>) -> Self {
233 Self {
234 inner,
235 last_primary_key: Vec::new(),
236 last_match: None,
237 }
238 }
239}
240
241impl PrimaryKeyFilter for CachedPrimaryKeyFilter {
242 fn matches(&mut self, pk: &[u8]) -> mito_codec::error::Result<bool> {
243 if let Some(last_match) = self.last_match
244 && self.last_primary_key == pk
245 {
246 return Ok(last_match);
247 }
248
249 let matched = self.inner.matches(pk)?;
250 self.last_primary_key.clear();
251 self.last_primary_key.extend_from_slice(pk);
252 self.last_match = Some(matched);
253 Ok(matched)
254 }
255}
256
257pub(crate) struct BulkFilterPlan {
263 pub(crate) remaining_simple_filters: Vec<SimpleFilterContext>,
265 pub(crate) pk_filters: Option<Arc<Vec<SimpleFilterEvaluator>>>,
268}
269
270pub(crate) fn build_primary_key_filter(
276 sst_metadata: &RegionMetadataRef,
277 expected_metadata: Option<&RegionMetadata>,
278 predicate: Option<&Predicate>,
279) -> Option<CachedPrimaryKeyFilter> {
280 let filters = simple_tag_filters(sst_metadata, expected_metadata, predicate)
281 .into_iter()
282 .map(|(_, filter)| filter)
283 .collect::<Vec<_>>();
284 if filters.is_empty() {
285 return None;
286 }
287
288 let codec = build_primary_key_codec(sst_metadata.as_ref());
289 let filter = codec.primary_key_filter(sst_metadata, Arc::new(filters));
290 Some(CachedPrimaryKeyFilter::new(filter))
291}
292
293pub(crate) fn simple_tag_filters(
295 sst_metadata: &RegionMetadataRef,
296 expected_metadata: Option<&RegionMetadata>,
297 predicate: Option<&Predicate>,
298) -> Vec<(Expr, SimpleFilterEvaluator)> {
299 predicate
300 .into_iter()
301 .flat_map(|predicate| predicate.exprs())
302 .filter_map(|expr| {
303 SimpleFilterContext::new_opt(sst_metadata, expected_metadata, expr)
304 .map(|filter_ctx| (expr, filter_ctx))
305 })
306 .filter_map(|(expr, filter_ctx)| {
307 (filter_ctx.semantic_type() == SemanticType::Tag)
308 .then(|| {
309 filter_ctx
310 .filter()
311 .as_filter()
312 .cloned()
313 .map(|filter| (expr.clone(), filter))
314 })
315 .flatten()
316 })
317 .collect()
318}
319
320pub(crate) struct ReaderFilterPlan {
328 pub(crate) remaining_simple_filters: Vec<SimpleFilterContext>,
333 pub(crate) prefilter_builder: Option<PrefilterContextBuilder>,
337}
338
339pub(crate) fn build_bulk_filter_plan(
340 read_format: &FlatReadFormat,
341 predicate: Option<&Predicate>,
342) -> BulkFilterPlan {
343 let metadata = read_format.metadata();
344 let simple_filters: Vec<SimpleFilterContext> = predicate
347 .into_iter()
348 .flat_map(|predicate| {
349 predicate
350 .exprs()
351 .iter()
352 .filter_map(|expr| SimpleFilterContext::new_opt(metadata, None, expr))
353 })
354 .collect();
355
356 if read_format.batch_has_raw_pk_columns() || metadata.primary_key.is_empty() {
360 return BulkFilterPlan {
361 remaining_simple_filters: simple_filters,
362 pk_filters: None,
363 };
364 }
365
366 let mut remaining_simple_filters = Vec::new();
367 let mut pk_filters = Vec::new();
368
369 for filter_ctx in simple_filters {
370 let pk_filter = filter_ctx.filter().as_filter().and_then(|filter| {
373 (filter_ctx.semantic_type() == SemanticType::Tag).then(|| filter.clone())
374 });
375
376 if let Some(pk_filter) = pk_filter {
377 pk_filters.push(pk_filter);
378 } else {
379 remaining_simple_filters.push(filter_ctx);
380 }
381 }
382
383 BulkFilterPlan {
384 remaining_simple_filters,
385 pk_filters: (!pk_filters.is_empty()).then_some(Arc::new(pk_filters)),
386 }
387}
388
389#[allow(clippy::too_many_arguments)]
406pub(crate) fn build_reader_filter_plan(
407 predicate: Option<&Predicate>,
408 expected_metadata: Option<&RegionMetadata>,
409 pre_filter_mode: PreFilterMode,
410 enable_predicate_prefilter: bool,
411 postpone_time_index_filter: bool,
412 read_format: &FlatReadFormat,
413 codec: &Arc<dyn PrimaryKeyCodec>,
414 parquet_metadata: &ParquetMetaData,
415) -> ReaderFilterPlan {
416 let Some(predicate) = predicate else {
417 return ReaderFilterPlan {
418 remaining_simple_filters: Vec::new(),
419 prefilter_builder: None,
420 };
421 };
422
423 let metadata = read_format.metadata();
424 let mut prefilter_simple_filters = Vec::new();
425 let mut remaining_simple_filters = Vec::new();
426 let mut prefilter_physical_filters = Vec::new();
427 let mut primary_key_filters = Vec::new();
428 let mut pk_filter_contexts = Vec::new();
429
430 let field_prefilter_enabled = pre_filter_mode == PreFilterMode::All;
434 let need_pk_prefilter = !read_format.batch_has_raw_pk_columns();
438
439 let can_direct_prefilter = |semantic_type: SemanticType| -> bool {
442 match semantic_type {
443 SemanticType::Tag => !need_pk_prefilter,
444 SemanticType::Field => field_prefilter_enabled,
445 SemanticType::Timestamp => true,
446 }
447 };
448
449 for expr in predicate.exprs() {
450 if let Some(filter_ctx) = SimpleFilterContext::new_opt(metadata, expected_metadata, expr) {
453 if !enable_predicate_prefilter {
454 remaining_simple_filters.push(filter_ctx);
455 continue;
456 }
457
458 let Some(filter) = filter_ctx.filter().as_filter() else {
462 remaining_simple_filters.push(filter_ctx);
463 continue;
464 };
465
466 if postpone_time_index_filter && filter_ctx.semantic_type() == SemanticType::Timestamp {
467 remaining_simple_filters.push(filter_ctx);
468 continue;
469 }
470
471 let direct_prefilter = can_direct_prefilter(filter_ctx.semantic_type());
474 if direct_prefilter {
475 assert!(
476 read_format
477 .arrow_schema()
478 .column_with_name(filter.column_name())
479 .is_some(),
480 "Column '{}' is not present in the arrow schema {:?}",
481 filter.column_name(),
482 read_format.arrow_schema(),
483 );
484 prefilter_simple_filters.push(filter_ctx);
485 continue;
486 }
487
488 if need_pk_prefilter && filter_ctx.semantic_type() == SemanticType::Tag {
490 primary_key_filters.push(filter.clone());
491 pk_filter_contexts.push(filter_ctx);
492 } else {
493 remaining_simple_filters.push(filter_ctx);
494 }
495 continue;
496 }
497
498 if !enable_predicate_prefilter {
499 continue;
500 }
501
502 if let Some(filter) =
507 PhysicalFilterContext::new_opt(metadata, expected_metadata, read_format, expr)
508 && can_direct_prefilter(filter.semantic_type())
509 {
510 prefilter_physical_filters.push(filter);
511 }
512 }
513
514 if !enable_predicate_prefilter {
515 return ReaderFilterPlan {
516 remaining_simple_filters,
517 prefilter_builder: None,
518 };
519 }
520
521 let pk_filter_expr_strs = (!pk_filter_contexts.is_empty()).then(|| {
522 let mut expr_strs = pk_filter_contexts
523 .iter()
524 .map(|filter_ctx| filter_ctx.expr_str().to_string())
525 .collect::<Vec<_>>();
526 expr_strs.sort();
527 SmallVec::from_vec(expr_strs)
528 });
529 let pk_filter_exprs =
530 (!primary_key_filters.is_empty()).then_some(Arc::new(primary_key_filters));
531 let schema_version = expected_metadata
532 .map(|metadata| metadata.schema_version)
533 .unwrap_or_else(|| read_format.metadata().schema_version);
534 let prefilter_builder = PrefilterContextBuilder::new(
535 read_format,
536 codec,
537 pk_filter_exprs,
538 pk_filter_expr_strs,
539 prefilter_simple_filters.clone(),
540 prefilter_physical_filters,
541 schema_version,
542 parquet_metadata,
543 );
544
545 if prefilter_builder.is_some() {
546 ReaderFilterPlan {
547 remaining_simple_filters,
548 prefilter_builder,
549 }
550 } else {
551 remaining_simple_filters.extend(prefilter_simple_filters);
554 remaining_simple_filters.extend(pk_filter_contexts);
555 ReaderFilterPlan {
556 remaining_simple_filters,
557 prefilter_builder: None,
558 }
559 }
560}
561
562pub(crate) struct PrefilterContext {
564 pk_filter: Option<Box<dyn PrimaryKeyFilter>>,
566 filters: Vec<SimpleFilterContext>,
568 physical_filters: Vec<PhysicalFilterContext>,
571 schema_version: u64,
573 pk_filter_expr_strs: Option<SmallVec<[String; 1]>>,
575 arrow_schema: SchemaRef,
577 proven_simple_filters: Vec<bool>,
579}
580
581pub(crate) struct PrefilterContextBuilder {
587 pk_filters: Option<Arc<Vec<SimpleFilterEvaluator>>>,
588 pk_filter_expr_strs: Option<SmallVec<[String; 1]>>,
589 filters: Vec<SimpleFilterContext>,
590 physical_filters: Vec<PhysicalFilterContext>,
591 codec: Arc<dyn PrimaryKeyCodec>,
592 metadata: RegionMetadataRef,
593 schema_version: u64,
594 arrow_schema: SchemaRef,
595 proven_simple_filters: Vec<Vec<bool>>,
597}
598
599impl PrefilterContextBuilder {
600 #[allow(clippy::too_many_arguments)]
607 pub(crate) fn new(
608 read_format: &FlatReadFormat,
609 codec: &Arc<dyn PrimaryKeyCodec>,
610 primary_key_filters: Option<Arc<Vec<SimpleFilterEvaluator>>>,
611 primary_key_filter_expr_strs: Option<SmallVec<[String; 1]>>,
612 filters: Vec<SimpleFilterContext>,
613 physical_filters: Vec<PhysicalFilterContext>,
614 schema_version: u64,
615 parquet_metadata: &ParquetMetaData,
616 ) -> Option<Self> {
617 let metadata = read_format.metadata();
618 let use_raw_tag_columns = read_format.batch_has_raw_pk_columns();
619 let pk_filters = (!use_raw_tag_columns)
620 .then_some(primary_key_filters)
621 .flatten()
622 .filter(|filters| !filters.is_empty());
623 let pk_filter_expr_strs = pk_filters
624 .is_some()
625 .then_some(primary_key_filter_expr_strs)
626 .flatten();
627
628 let mut prefilter_column_names = HashSet::new();
629 for filter_ctx in &filters {
630 if let MaybeFilter::Filter(filter) = filter_ctx.filter() {
631 prefilter_column_names.insert(filter.column_name().to_string());
632 }
633 }
634
635 if pk_filters.is_some() {
636 prefilter_column_names.insert(PRIMARY_KEY_COLUMN_NAME.to_string());
637 }
638
639 for filter_ctx in &physical_filters {
640 prefilter_column_names.insert(filter_ctx.column_name().to_string());
641 }
642
643 let prefilter_count =
644 compute_projection_count(&prefilter_column_names, read_format.arrow_schema());
645
646 if prefilter_count == 0 {
647 return None;
648 }
649
650 let total_count = read_format.parquet_read_columns().root_indices().len();
651 let remaining_count = total_count.saturating_sub(prefilter_count);
652 if pk_filters.is_none() && prefilter_count >= total_count {
653 return None;
654 }
655
656 if pk_filters.is_none()
657 && !should_use_prefilter(prefilter_count, remaining_count, total_count)
658 {
659 return None;
660 }
661
662 let proven_simple_filters =
663 simple_filter_stats_proofs(read_format, parquet_metadata.row_groups(), &filters);
664
665 Some(Self {
666 pk_filters,
667 pk_filter_expr_strs,
668 filters,
669 physical_filters,
670 codec: Arc::clone(codec),
671 metadata: metadata.clone(),
672 schema_version,
673 arrow_schema: read_format.arrow_schema().clone(),
674 proven_simple_filters,
675 })
676 }
677
678 pub(crate) fn build(&self, row_group_idx: usize) -> PrefilterContext {
680 let pk_filter = self
681 .build_primary_key_filter()
682 .map(|filter| Box::new(filter) as Box<dyn PrimaryKeyFilter>);
683 PrefilterContext {
684 pk_filter,
685 filters: self.filters.clone(),
686 physical_filters: self.physical_filters.clone(),
687 schema_version: self.schema_version,
688 pk_filter_expr_strs: self.pk_filter_expr_strs.clone(),
689 arrow_schema: self.arrow_schema.clone(),
690 proven_simple_filters: self
691 .proven_simple_filters
692 .get(row_group_idx)
693 .cloned()
694 .unwrap_or_else(|| vec![false; self.filters.len()]),
695 }
696 }
697
698 pub(crate) fn build_primary_key_filter(&self) -> Option<CachedPrimaryKeyFilter> {
700 self.pk_filters.as_ref().map(|pk_filters| {
701 let filter = self
702 .codec
703 .primary_key_filter(&self.metadata, Arc::clone(pk_filters));
704 CachedPrimaryKeyFilter::new(filter)
705 })
706 }
707}
708
709const PREFILTER_COLUMN_RATIO_THRESHOLD: f64 = 0.5;
710const PREFILTER_MIN_REMAINING_COLUMNS: usize = 2;
711
712fn simple_filter_stats_proofs(
715 read_format: &FlatReadFormat,
716 row_groups: &[parquet::file::metadata::RowGroupMetaData],
717 filters: &[SimpleFilterContext],
718) -> Vec<Vec<bool>> {
719 let mut proofs = vec![vec![false; filters.len()]; row_groups.len()];
720 for (filter_idx, filter_ctx) in filters.iter().enumerate() {
721 let Some((filter, literal)) = eligible_simple_filter(read_format, filter_ctx) else {
722 continue;
723 };
724 let (StatValues::Values(mins), StatValues::Values(maxs), StatValues::Values(null_counts)) = (
725 read_format.min_values(row_groups, filter_ctx.column_id()),
726 read_format.max_values(row_groups, filter_ctx.column_id()),
727 read_format.null_counts(row_groups, filter_ctx.column_id()),
728 ) else {
729 continue;
730 };
731 for (row_group_idx, proof) in proofs.iter_mut().enumerate() {
732 proof[filter_idx] = simple_filter_is_true_by_values(
733 filter,
734 &literal,
735 stat_value_at(&mins, row_group_idx),
736 stat_value_at(&maxs, row_group_idx),
737 stat_value_at(&null_counts, row_group_idx),
738 );
739 }
740 }
741 proofs
742}
743
744fn eligible_simple_filter<'a>(
745 read_format: &FlatReadFormat,
746 filter_ctx: &'a SimpleFilterContext,
747) -> Option<(&'a SimpleFilterEvaluator, Value)> {
748 if filter_ctx.semantic_type() != SemanticType::Field {
749 return None;
750 }
751 let filter = filter_ctx.filter().as_filter()?;
752 let literal = filter.literal_value()?;
753 let column = read_format
754 .metadata()
755 .column_by_id(filter_ctx.column_id())?;
756 column_type_matches_literal(&column.column_schema.data_type, &literal)
757 .then_some((filter, literal))
758}
759
760fn simple_filter_is_true_by_values(
761 filter: &SimpleFilterEvaluator,
762 literal: &Value,
763 min: Option<Value>,
764 max: Option<Value>,
765 null_count: Option<Value>,
766) -> bool {
767 let (Some(min), Some(max), Some(null_count)) = (min, max, null_count) else {
768 return false;
769 };
770 if null_count != Value::UInt64(0)
771 || !same_supported_value_type(&min, literal)
772 || !same_supported_value_type(&max, literal)
773 || min > max
774 {
775 return false;
776 }
777
778 if filter.is_gt() {
779 min > *literal
780 } else if filter.is_gt_eq() {
781 min >= *literal
782 } else if filter.is_lt() {
783 max < *literal
784 } else if filter.is_lt_eq() {
785 max <= *literal
786 } else if filter.is_eq() {
787 min == *literal && max == *literal
788 } else if filter.is_not_eq() {
789 max < *literal || min > *literal
790 } else {
791 false
792 }
793}
794
795fn column_type_matches_literal(data_type: &ConcreteDataType, literal: &Value) -> bool {
796 matches!(
797 (data_type, literal),
798 (ConcreteDataType::Int32(_), Value::Int32(_))
799 | (ConcreteDataType::UInt32(_), Value::UInt32(_))
800 | (ConcreteDataType::Int64(_), Value::Int64(_))
801 | (ConcreteDataType::UInt64(_), Value::UInt64(_))
802 )
803}
804
805fn same_supported_value_type(left: &Value, right: &Value) -> bool {
806 matches!(
807 (left, right),
808 (Value::Int32(_), Value::Int32(_))
809 | (Value::UInt32(_), Value::UInt32(_))
810 | (Value::Int64(_), Value::Int64(_))
811 | (Value::UInt64(_), Value::UInt64(_))
812 )
813}
814
815fn stat_value_at(values: &datatypes::arrow::array::ArrayRef, index: usize) -> Option<Value> {
816 let scalar = ScalarValue::try_from_array(values, index).ok()?;
817 Value::try_from(scalar).ok()
818}
819
820pub(crate) struct PrefilterResult {
822 pub(crate) refined_selection: RowSelection,
824 pub(crate) filtered_rows: usize,
826}
827
828fn compute_projection_mask(
833 column_names: &HashSet<String>,
834 arrow_schema: &datatypes::arrow::datatypes::SchemaRef,
835 parquet_schema: &SchemaDescriptor,
836) -> ProjectionMask {
837 ProjectionMask::roots(
838 parquet_schema,
839 projection_indices(column_names, arrow_schema),
840 )
841}
842
843fn compute_projection_count(
844 column_names: &HashSet<String>,
845 arrow_schema: &datatypes::arrow::datatypes::SchemaRef,
846) -> usize {
847 projection_indices(column_names, arrow_schema).len()
848}
849
850fn projection_indices(
851 column_names: &HashSet<String>,
852 arrow_schema: &datatypes::arrow::datatypes::SchemaRef,
853) -> Vec<usize> {
854 let mut projection_indices: Vec<usize> = column_names
855 .iter()
856 .filter_map(|name| arrow_schema.column_with_name(name).map(|(index, _)| index))
857 .collect();
858 projection_indices.sort_unstable();
859 projection_indices.dedup();
860 projection_indices
861}
862
863fn should_use_prefilter(
864 prefilter_count: usize,
865 remaining_count: usize,
866 total_count: usize,
867) -> bool {
868 if remaining_count == 0 {
869 return false;
870 }
871
872 if remaining_count < PREFILTER_MIN_REMAINING_COLUMNS {
873 return false;
874 }
875
876 let ratio = prefilter_count as f64 / total_count as f64;
877 ratio <= PREFILTER_COLUMN_RATIO_THRESHOLD
878}
879
880pub(crate) async fn execute_prefilter(
881 prefilter_ctx: &mut PrefilterContext,
882 reader_builder: &RowGroupReaderBuilder,
883 build_ctx: &RowGroupBuildContext<'_>,
884) -> Result<PrefilterResult> {
885 let entries = build_prefilter_cache_entries(prefilter_ctx, reader_builder, build_ctx);
886
887 if entries.is_empty() {
888 return execute_prefilter_by_reading_columns(prefilter_ctx, reader_builder, build_ctx)
889 .await;
890 }
891
892 execute_prefilter_with_result_cache(prefilter_ctx, reader_builder, build_ctx, entries).await
893}
894
895async fn execute_prefilter_with_result_cache(
896 prefilter_ctx: &mut PrefilterContext,
897 reader_builder: &RowGroupReaderBuilder,
898 build_ctx: &RowGroupBuildContext<'_>,
899 entries: Vec<PrefilterEntry>,
900) -> Result<PrefilterResult> {
901 let non_cacheable_physical = non_cacheable_physical_filters(prefilter_ctx);
902 let mut hit_mask: Option<BooleanBuffer> = None;
903 let mut misses = Vec::new();
904 for entry in entries {
905 let Some(key) = &entry.key else {
906 misses.push(entry);
907 continue;
908 };
909
910 if let Some(mask) = reader_builder.cache_strategy().get_prefilter_result(key) {
911 hit_mask = Some(match hit_mask {
912 Some(hit_mask) => hit_mask.bitand(mask.as_ref()),
913 None => mask.as_ref().clone(),
914 });
915 } else {
916 misses.push(entry);
917 }
918 }
919
920 if misses.is_empty() && non_cacheable_physical.is_empty() {
921 let combined_mask = hit_mask.unwrap_or_else(|| BooleanBuffer::new_set(0));
922 let refined_selection =
923 refined_selection_from_mask(&combined_mask, &build_ctx.row_selection);
924 let rows_before_filter = rows_before_filter(reader_builder, build_ctx);
925 let filtered_rows = rows_before_filter.saturating_sub(refined_selection.row_count());
926 return Ok(PrefilterResult {
927 refined_selection,
928 filtered_rows,
929 });
930 }
931
932 let mut uncached_entries = misses;
933 uncached_entries.extend(
934 non_cacheable_physical
935 .iter()
936 .copied()
937 .map(|idx| PrefilterEntry::without_cache(PrefilterEntryKind::Physical(idx))),
938 );
939 let (uncached_mask, read_rows) =
940 build_prefilter_masks(prefilter_ctx, reader_builder, build_ctx, &uncached_entries).await?;
941
942 let final_mask = match (hit_mask, uncached_mask) {
943 (Some(hit_mask), Some(uncached_mask)) => hit_mask.bitand(&uncached_mask),
944 (Some(hit_mask), None) => hit_mask,
945 (None, Some(uncached_mask)) => uncached_mask,
946 (None, None) => BooleanBuffer::new_set(read_rows),
947 };
948 debug_assert_eq!(final_mask.len(), read_rows);
949 let rows_selected = final_mask.count_set_bits();
950 let filtered_rows = read_rows.saturating_sub(rows_selected);
951 let refined_selection = refined_selection_from_mask(&final_mask, &build_ctx.row_selection);
952
953 Ok(PrefilterResult {
954 refined_selection,
955 filtered_rows,
956 })
957}
958
959fn non_cacheable_physical_filters(prefilter_ctx: &PrefilterContext) -> Vec<usize> {
960 prefilter_ctx
961 .physical_filters
962 .iter()
963 .enumerate()
964 .filter_map(|(idx, filter)| (!filter.is_immutable()).then_some(idx))
965 .collect()
966}
967
968async fn build_prefilter_masks(
969 prefilter_ctx: &mut PrefilterContext,
970 reader_builder: &RowGroupReaderBuilder,
971 build_ctx: &RowGroupBuildContext<'_>,
972 entries: &[PrefilterEntry],
973) -> Result<(Option<BooleanBuffer>, usize)> {
974 let prefilter_column_names = prefilter_column_names_for_entries(prefilter_ctx, entries);
975 let parquet_schema = reader_builder
976 .parquet_metadata()
977 .file_metadata()
978 .schema_descr();
979 let projection = compute_projection_mask(
980 &prefilter_column_names,
981 &prefilter_ctx.arrow_schema,
982 parquet_schema,
983 );
984
985 let mut stream = reader_builder
986 .build_with_projection(
987 build_ctx.row_group_idx,
988 build_ctx.row_selection.clone(),
989 projection,
990 build_ctx.fetch_metrics,
991 )
992 .await?;
993
994 let mut cache_builders = entries
995 .iter()
996 .map(|entry| entry.key.is_some().then(|| BooleanBufferBuilder::new(0)))
997 .collect::<Vec<_>>();
998 let mut combined_builder = (!entries.is_empty()).then(|| BooleanBufferBuilder::new(0));
999 let mut rows_before_filter = 0usize;
1000
1001 while let Some(batch_result) = stream.next().await {
1002 let batch = batch_result?;
1003 let num_rows = batch.num_rows();
1004 if num_rows == 0 {
1005 continue;
1006 }
1007 rows_before_filter += num_rows;
1008
1009 let mut batch_mask = BooleanBuffer::new_set(num_rows);
1010 for (idx, entry) in entries.iter().enumerate() {
1011 let mask = eval_entry_mask(
1012 &batch,
1013 prefilter_ctx,
1014 entry.kind,
1015 reader_builder.file_path(),
1016 )?;
1017 batch_mask = batch_mask.bitand(&mask);
1018 if let Some(Some(builder)) = cache_builders.get_mut(idx) {
1019 builder.append_buffer(&mask);
1020 }
1021 }
1022 if let Some(builder) = &mut combined_builder {
1023 builder.append_buffer(&batch_mask);
1024 }
1025 }
1026
1027 for (entry, builder) in entries.iter().zip(cache_builders) {
1028 if let (Some(key), Some(mut builder)) = (&entry.key, builder) {
1029 reader_builder
1030 .cache_strategy()
1031 .put_prefilter_result(key.clone(), Arc::new(builder.finish()));
1032 }
1033 }
1034
1035 Ok((
1036 combined_builder.map(|mut builder| builder.finish()),
1037 rows_before_filter,
1038 ))
1039}
1040
1041fn prefilter_column_names_for_entries(
1042 prefilter_ctx: &PrefilterContext,
1043 entries: &[PrefilterEntry],
1044) -> HashSet<String> {
1045 let mut prefilter_column_names = HashSet::new();
1046 for entry in entries {
1047 match entry.kind {
1048 PrefilterEntryKind::Simple(idx) => {
1049 if let MaybeFilter::Filter(filter) = prefilter_ctx.filters[idx].filter() {
1050 prefilter_column_names.insert(filter.column_name().to_string());
1051 }
1052 }
1053 PrefilterEntryKind::Physical(idx) => {
1054 prefilter_column_names.insert(
1055 prefilter_ctx.physical_filters[idx]
1056 .column_name()
1057 .to_string(),
1058 );
1059 }
1060 PrefilterEntryKind::PkGroup => {
1061 prefilter_column_names.insert(PRIMARY_KEY_COLUMN_NAME.to_string());
1062 }
1063 }
1064 }
1065 prefilter_column_names
1066}
1067
1068async fn execute_prefilter_by_reading_columns(
1069 prefilter_ctx: &mut PrefilterContext,
1070 reader_builder: &RowGroupReaderBuilder,
1071 build_ctx: &RowGroupBuildContext<'_>,
1072) -> Result<PrefilterResult> {
1073 let entries = all_prefilter_entries(prefilter_ctx);
1074 if entries.is_empty() {
1075 return Ok(identity_prefilter_result(reader_builder, build_ctx));
1076 }
1077 let (mask, rows_before_filter) =
1078 build_prefilter_masks(prefilter_ctx, reader_builder, build_ctx, &entries).await?;
1079
1080 let final_mask = mask.unwrap_or_else(|| BooleanBuffer::new_set(rows_before_filter));
1081 let rows_selected = final_mask.count_set_bits();
1082 let filtered_rows = rows_before_filter.saturating_sub(rows_selected);
1083 let refined_selection = refined_selection_from_mask(&final_mask, &build_ctx.row_selection);
1084
1085 Ok(PrefilterResult {
1086 refined_selection,
1087 filtered_rows,
1088 })
1089}
1090
1091fn all_prefilter_entries(prefilter_ctx: &PrefilterContext) -> Vec<PrefilterEntry> {
1092 let mut entries = Vec::new();
1093 if prefilter_ctx.pk_filter.is_some() {
1094 entries.push(PrefilterEntry::without_cache(PrefilterEntryKind::PkGroup));
1095 }
1096 entries.extend(
1097 prefilter_ctx
1098 .filters
1099 .iter()
1100 .enumerate()
1101 .filter(|(idx, _)| {
1102 !prefilter_ctx
1103 .proven_simple_filters
1104 .get(*idx)
1105 .copied()
1106 .unwrap_or(false)
1107 })
1108 .map(|(idx, _)| PrefilterEntry::without_cache(PrefilterEntryKind::Simple(idx))),
1109 );
1110 entries.extend(
1111 prefilter_ctx
1112 .physical_filters
1113 .iter()
1114 .enumerate()
1115 .map(|(idx, _)| PrefilterEntry::without_cache(PrefilterEntryKind::Physical(idx))),
1116 );
1117 entries
1118}
1119
1120#[derive(Clone, Copy)]
1121enum PrefilterEntryKind {
1122 Simple(usize),
1123 Physical(usize),
1124 PkGroup,
1125}
1126
1127struct PrefilterEntry {
1128 kind: PrefilterEntryKind,
1129 key: Option<PrefilterKey>,
1130}
1131
1132impl PrefilterEntry {
1133 fn without_cache(kind: PrefilterEntryKind) -> Self {
1134 Self { kind, key: None }
1135 }
1136}
1137
1138fn build_prefilter_cache_entries(
1139 prefilter_ctx: &PrefilterContext,
1140 reader_builder: &RowGroupReaderBuilder,
1141 build_ctx: &RowGroupBuildContext<'_>,
1142) -> Vec<PrefilterEntry> {
1143 let row_selection = PrefilterKey::row_selection_snapshot(build_ctx.row_selection.as_ref());
1144 let file_id = reader_builder.file_handle().file_id().file_id();
1145 let row_group_idx = build_ctx.row_group_idx as u32;
1146 let mut entries = Vec::new();
1147
1148 for (idx, filter_ctx) in prefilter_ctx.filters.iter().enumerate() {
1149 if prefilter_ctx
1150 .proven_simple_filters
1151 .get(idx)
1152 .copied()
1153 .unwrap_or(false)
1154 {
1155 continue;
1156 }
1157 entries.push(PrefilterEntry {
1158 kind: PrefilterEntryKind::Simple(idx),
1159 key: Some(PrefilterKey::new(
1160 file_id,
1161 row_group_idx,
1162 row_selection.clone(),
1163 prefilter_ctx.schema_version,
1164 smallvec![filter_ctx.expr_str().to_string()],
1165 )),
1166 });
1167 }
1168
1169 for (idx, filter_ctx) in prefilter_ctx.physical_filters.iter().enumerate() {
1170 if !filter_ctx.is_immutable() {
1171 continue;
1172 }
1173 entries.push(PrefilterEntry {
1174 kind: PrefilterEntryKind::Physical(idx),
1175 key: Some(PrefilterKey::new(
1176 file_id,
1177 row_group_idx,
1178 row_selection.clone(),
1179 prefilter_ctx.schema_version,
1180 smallvec![filter_ctx.expr_str().to_string()],
1181 )),
1182 });
1183 }
1184
1185 if prefilter_ctx.pk_filter.is_some()
1186 && let Some(exprs) = &prefilter_ctx.pk_filter_expr_strs
1187 {
1188 entries.push(PrefilterEntry {
1189 kind: PrefilterEntryKind::PkGroup,
1190 key: Some(PrefilterKey::new(
1191 file_id,
1192 row_group_idx,
1193 row_selection,
1194 prefilter_ctx.schema_version,
1195 exprs.clone(),
1196 )),
1197 });
1198 }
1199
1200 entries
1201}
1202
1203fn identity_prefilter_result(
1204 reader_builder: &RowGroupReaderBuilder,
1205 build_ctx: &RowGroupBuildContext<'_>,
1206) -> PrefilterResult {
1207 let row_count = reader_builder
1208 .parquet_metadata()
1209 .row_group(build_ctx.row_group_idx)
1210 .num_rows() as usize;
1211 PrefilterResult {
1212 refined_selection: identity_row_selection(&build_ctx.row_selection, row_count),
1213 filtered_rows: 0,
1214 }
1215}
1216
1217fn identity_row_selection(
1218 original_selection: &Option<RowSelection>,
1219 row_count: usize,
1220) -> RowSelection {
1221 original_selection
1222 .clone()
1223 .unwrap_or_else(|| RowSelection::from(vec![RowSelector::select(row_count)]))
1224}
1225
1226fn rows_before_filter(
1227 reader_builder: &RowGroupReaderBuilder,
1228 build_ctx: &RowGroupBuildContext<'_>,
1229) -> usize {
1230 build_ctx.row_selection.as_ref().map_or_else(
1231 || {
1232 reader_builder
1233 .parquet_metadata()
1234 .row_group(build_ctx.row_group_idx)
1235 .num_rows() as usize
1236 },
1237 RowSelection::row_count,
1238 )
1239}
1240
1241fn refined_selection_from_mask(
1242 mask: &BooleanBuffer,
1243 original_selection: &Option<RowSelection>,
1244) -> RowSelection {
1245 if mask.is_empty() || mask.count_set_bits() == 0 {
1246 return RowSelection::from(vec![]);
1247 }
1248
1249 let prefilter_selection = RowSelection::from_filters(&[BooleanArray::from(mask.clone())]);
1250 match original_selection {
1251 Some(original) => original.and_then(&prefilter_selection),
1252 None => prefilter_selection,
1253 }
1254}
1255
1256fn eval_entry_mask(
1257 batch: &RecordBatch,
1258 prefilter_ctx: &mut PrefilterContext,
1259 kind: PrefilterEntryKind,
1260 file_path: &str,
1261) -> Result<BooleanBuffer> {
1262 match kind {
1263 PrefilterEntryKind::Simple(idx) => {
1264 eval_simple_filter_mask(batch, &prefilter_ctx.filters[idx], file_path)
1265 }
1266 PrefilterEntryKind::Physical(idx) => {
1267 eval_physical_filter_mask(batch, &prefilter_ctx.physical_filters[idx], file_path)
1268 }
1269 PrefilterEntryKind::PkGroup => {
1270 let pk_filter = prefilter_ctx.pk_filter.as_mut().context(UnexpectedSnafu {
1271 reason: "Missing primary key filter for prefilter cache entry",
1272 })?;
1273 primary_key_filter_mask(batch, pk_filter.as_mut())
1274 }
1275 }
1276}
1277
1278pub(crate) fn primary_key_filter_mask(
1280 batch: &RecordBatch,
1281 pk_filter: &mut dyn PrimaryKeyFilter,
1282) -> Result<BooleanBuffer> {
1283 let (pk_column_index, _) = batch
1284 .schema()
1285 .column_with_name(PRIMARY_KEY_COLUMN_NAME)
1286 .context(UnexpectedSnafu {
1287 reason: "Primary key column not found in prefilter batch",
1288 })?;
1289 let matched_row_ranges = matching_row_ranges_by_primary_key(batch, pk_column_index, pk_filter)?;
1290 let mut builder = BooleanBufferBuilder::new(batch.num_rows());
1291 builder.append_n(batch.num_rows(), false);
1292 for range in matched_row_ranges {
1293 for row in range {
1294 builder.set_bit(row, true);
1295 }
1296 }
1297 Ok(builder.finish())
1298}
1299
1300fn eval_simple_filter_mask(
1301 batch: &RecordBatch,
1302 filter_ctx: &SimpleFilterContext,
1303 file_path: &str,
1304) -> Result<BooleanBuffer> {
1305 let filter = match filter_ctx.filter() {
1306 MaybeFilter::Filter(filter) => filter,
1307 MaybeFilter::Matched => return Ok(BooleanBuffer::new_set(batch.num_rows())),
1308 MaybeFilter::Pruned => return Ok(BooleanBuffer::new_unset(batch.num_rows())),
1309 };
1310
1311 let (idx, _) = batch
1312 .schema()
1313 .column_with_name(filter.column_name())
1314 .with_context(|| UnexpectedSnafu {
1315 reason: format!(
1316 "Prefilter column '{}' (id {}) not found in batch for file {}",
1317 filter.column_name(),
1318 filter_ctx.column_id(),
1319 file_path
1320 ),
1321 })?;
1322 let column = batch.column(idx).clone();
1323 filter.evaluate_array(&column).context(RecordBatchSnafu)
1324}
1325
1326fn eval_physical_filter_mask(
1327 batch: &RecordBatch,
1328 filter_ctx: &PhysicalFilterContext,
1329 file_path: &str,
1330) -> Result<BooleanBuffer> {
1331 let filter = filter_ctx.filter();
1332
1333 let (idx, _) = batch
1334 .schema()
1335 .column_with_name(filter_ctx.column_name())
1336 .with_context(|| UnexpectedSnafu {
1337 reason: format!(
1338 "Prefilter physical column '{}' (id {}) not found in batch for file {}",
1339 filter_ctx.column_name(),
1340 filter_ctx.column_id(),
1341 file_path
1342 ),
1343 })?;
1344 let column = batch.column(idx).clone();
1345
1346 let record_batch = RecordBatch::try_new(filter_ctx.schema().clone(), vec![column])
1347 .context(NewRecordBatchSnafu)?;
1348 let evaluated = filter
1349 .evaluate(&record_batch)
1350 .context(EvalPartitionFilterSnafu)?;
1351 let array = evaluated
1352 .into_array(record_batch.num_rows())
1353 .context(EvalPartitionFilterSnafu)?;
1354 let boolean_array = array
1355 .as_any()
1356 .downcast_ref::<BooleanArray>()
1357 .context(UnexpectedSnafu {
1358 reason: "Failed to downcast physical filter result to BooleanArray",
1359 })?;
1360 let mut result = boolean_array.values().clone();
1363 if let Some(nulls) = boolean_array.nulls() {
1364 result = result.bitand(nulls.inner());
1365 }
1366 Ok(result)
1367}
1368
1369#[cfg(test)]
1370mod tests {
1371 use std::sync::Arc;
1372 use std::sync::atomic::{AtomicUsize, Ordering};
1373
1374 use bytes::Bytes;
1375 use common_recordbatch::filter::SimpleFilterEvaluator;
1376 use datafusion_common::ScalarValue;
1377 use datafusion_expr::{col, lit};
1378 use datatypes::arrow::array::{
1379 ArrayRef, DictionaryArray, Int32Array, TimestampMillisecondArray, UInt8Array, UInt32Array,
1380 UInt64Array,
1381 };
1382 use datatypes::arrow::datatypes::{DataType, Field, Schema, UInt32Type};
1383 use datatypes::arrow::record_batch::RecordBatch;
1384 use datatypes::prelude::ConcreteDataType;
1385 use datatypes::value::Value;
1386 use mito_codec::row_converter::{PrimaryKeyFilter, build_primary_key_codec};
1387 use parquet::arrow::ArrowWriter;
1388 use parquet::arrow::arrow_reader::{ParquetRecordBatchReaderBuilder, RowSelector};
1389 use store_api::codec::PrimaryKeyEncoding;
1390 use store_api::metadata::RegionMetadataBuilder;
1391 use store_api::region_request::{AlterKind, ModifyColumnType};
1392
1393 use super::*;
1394 use crate::read::read_columns::ReadColumns;
1395 use crate::sst::internal_fields;
1396 use crate::sst::parquet::flat_format::{FlatReadFormat, primary_key_column_index};
1397 use crate::test_util::sst_util::{
1398 new_primary_key, new_record_batch_with_custom_sequence, sst_region_metadata,
1399 sst_region_metadata_with_encoding,
1400 };
1401
1402 struct CountingPrimaryKeyFilter {
1403 hits: Arc<AtomicUsize>,
1404 expected: Vec<u8>,
1405 }
1406
1407 impl PrimaryKeyFilter for CountingPrimaryKeyFilter {
1408 fn matches(&mut self, pk: &[u8]) -> mito_codec::error::Result<bool> {
1409 self.hits.fetch_add(1, Ordering::Relaxed);
1410 Ok(pk == self.expected.as_slice())
1411 }
1412 }
1413
1414 #[test]
1415 fn test_cached_primary_key_filter_reuses_previous_result() {
1416 let expected = new_primary_key(&["a", "x"]);
1417 let hits = Arc::new(AtomicUsize::new(0));
1418 let mut filter = CachedPrimaryKeyFilter::new(Box::new(CountingPrimaryKeyFilter {
1419 hits: Arc::clone(&hits),
1420 expected: expected.clone(),
1421 }));
1422
1423 assert!(filter.matches(expected.as_slice()).unwrap());
1424 assert!(filter.matches(expected.as_slice()).unwrap());
1425 assert!(
1426 !filter
1427 .matches(new_primary_key(&["b", "x"]).as_slice())
1428 .unwrap()
1429 );
1430
1431 assert_eq!(hits.load(Ordering::Relaxed), 2);
1432 }
1433
1434 fn new_test_filters(exprs: &[datafusion_expr::Expr]) -> Vec<SimpleFilterEvaluator> {
1435 exprs
1436 .iter()
1437 .filter_map(SimpleFilterEvaluator::try_new)
1438 .collect()
1439 }
1440
1441 fn new_simple_filter_contexts(
1442 metadata: &RegionMetadataRef,
1443 exprs: &[datafusion_expr::Expr],
1444 ) -> Vec<SimpleFilterContext> {
1445 exprs
1446 .iter()
1447 .filter_map(|expr| SimpleFilterContext::new_opt(metadata, None, expr))
1448 .collect()
1449 }
1450
1451 fn stats_metadata_with_options(
1452 values: &[Vec<u64>],
1453 nulls: Option<&[bool]>,
1454 writer_options: Option<parquet::file::properties::WriterProperties>,
1455 ) -> Arc<ParquetMetaData> {
1456 let first = new_record_batch_with_custom_sequence(&["a", "x"], 0, values[0].len(), 1);
1457 let mut bytes = Vec::new();
1458 let mut writer = ArrowWriter::try_new(&mut bytes, first.schema(), writer_options).unwrap();
1459 for (idx, values) in values.iter().enumerate() {
1460 let has_null = nulls.and_then(|nulls| nulls.get(idx)).copied() == Some(true);
1461 let batch = new_record_batch_with_custom_sequence(
1462 &["a", "x"],
1463 0,
1464 if has_null {
1465 values.len().max(2)
1466 } else {
1467 values.len()
1468 },
1469 1,
1470 );
1471 let mut columns = batch.columns().to_vec();
1472 columns[2] = Arc::new(if has_null {
1473 UInt64Array::from(vec![Some(values[0]), None])
1474 } else {
1475 UInt64Array::from_iter_values(values.iter().copied())
1476 });
1477 writer
1478 .write(&RecordBatch::try_new(batch.schema(), columns).unwrap())
1479 .unwrap();
1480 writer.flush().unwrap();
1481 }
1482 writer.close().unwrap();
1483 ParquetRecordBatchReaderBuilder::try_new(Bytes::from(bytes))
1484 .unwrap()
1485 .metadata()
1486 .clone()
1487 }
1488
1489 fn stats_metadata(values: &[Vec<u64>], nulls: Option<&[bool]>) -> Arc<ParquetMetaData> {
1490 stats_metadata_with_options(values, nulls, None)
1491 }
1492
1493 fn int32_stats_metadata(values: &[i32]) -> Arc<ParquetMetaData> {
1494 let batch = new_record_batch_with_custom_sequence(&["a", "x"], 0, values.len(), 1);
1495 let mut fields = batch
1496 .schema()
1497 .fields()
1498 .iter()
1499 .map(|field| field.as_ref().clone())
1500 .collect::<Vec<_>>();
1501 fields[2].set_data_type(DataType::Int32);
1502 let mut columns = batch.columns().to_vec();
1503 columns[2] = Arc::new(Int32Array::from(values.to_vec()));
1504 let batch = RecordBatch::try_new(Arc::new(Schema::new(fields)), columns).unwrap();
1505 let mut bytes = Vec::new();
1506 let mut writer = ArrowWriter::try_new(&mut bytes, batch.schema(), None).unwrap();
1507 writer.write(&batch).unwrap();
1508 writer.close().unwrap();
1509 ParquetRecordBatchReaderBuilder::try_new(Bytes::from(bytes))
1510 .unwrap()
1511 .metadata()
1512 .clone()
1513 }
1514
1515 fn metadata_with_field_type(data_type: ConcreteDataType) -> RegionMetadataRef {
1516 let mut builder = RegionMetadataBuilder::from_existing(sst_region_metadata());
1517 builder
1518 .alter(AlterKind::ModifyColumnTypes {
1519 columns: vec![ModifyColumnType {
1520 column_name: "field_0".to_string(),
1521 target_type: data_type,
1522 }],
1523 })
1524 .unwrap();
1525 Arc::new(builder.build().unwrap())
1526 }
1527
1528 fn new_physical_filter_contexts(
1529 metadata: &RegionMetadataRef,
1530 read_format: &FlatReadFormat,
1531 exprs: &[datafusion_expr::Expr],
1532 ) -> Vec<PhysicalFilterContext> {
1533 exprs
1534 .iter()
1535 .filter_map(|expr| PhysicalFilterContext::new_opt(metadata, None, read_format, expr))
1536 .collect()
1537 }
1538
1539 fn new_raw_batch(primary_keys: &[&[u8]], field_values: &[u64]) -> RecordBatch {
1540 assert_eq!(primary_keys.len(), field_values.len());
1541
1542 let metadata = Arc::new(sst_region_metadata());
1543 let arrow_schema = metadata.schema.arrow_schema();
1544 let field_column = arrow_schema
1545 .field(arrow_schema.index_of("field_0").unwrap())
1546 .clone();
1547 let time_index_column = arrow_schema
1548 .field(arrow_schema.index_of("ts").unwrap())
1549 .clone();
1550 let mut fields = vec![field_column, time_index_column];
1551 fields.extend(
1552 internal_fields()
1553 .into_iter()
1554 .map(|field| field.as_ref().clone()),
1555 );
1556 let schema = Arc::new(Schema::new(fields));
1557
1558 let mut dict_values = Vec::new();
1559 let mut keys = Vec::with_capacity(primary_keys.len());
1560 for pk in primary_keys {
1561 let key = dict_values
1562 .iter()
1563 .position(|existing: &&[u8]| existing == pk)
1564 .unwrap_or_else(|| {
1565 dict_values.push(*pk);
1566 dict_values.len() - 1
1567 });
1568 keys.push(key as u32);
1569 }
1570 let pk_array: ArrayRef = Arc::new(DictionaryArray::<UInt32Type>::new(
1571 UInt32Array::from(keys),
1572 Arc::new(BinaryArray::from_iter_values(dict_values.iter().copied())),
1573 ));
1574
1575 RecordBatch::try_new(
1576 schema,
1577 vec![
1578 Arc::new(UInt64Array::from(field_values.to_vec())),
1579 Arc::new(TimestampMillisecondArray::from_iter_values(
1580 0..primary_keys.len() as i64,
1581 )),
1582 pk_array,
1583 Arc::new(UInt64Array::from(vec![1; primary_keys.len()])),
1584 Arc::new(UInt8Array::from(vec![1; primary_keys.len()])),
1585 ],
1586 )
1587 .unwrap()
1588 }
1589
1590 fn new_prefilter_batch(primary_keys: &[&[u8]], field_values: &[u64]) -> RecordBatch {
1591 assert_eq!(primary_keys.len(), field_values.len());
1592
1593 let metadata = Arc::new(sst_region_metadata());
1594 let arrow_schema = metadata.schema.arrow_schema();
1595 let field_column = arrow_schema
1596 .field(arrow_schema.index_of("field_0").unwrap())
1597 .clone();
1598 let time_index_column = arrow_schema
1599 .field(arrow_schema.index_of("ts").unwrap())
1600 .clone();
1601 let schema = Arc::new(Schema::new(vec![
1602 field_column,
1603 time_index_column,
1604 internal_fields()[0].as_ref().clone(),
1605 ]));
1606
1607 let mut dict_values = Vec::new();
1608 let mut keys = Vec::with_capacity(primary_keys.len());
1609 for pk in primary_keys {
1610 let key = dict_values
1611 .iter()
1612 .position(|existing: &&[u8]| existing == pk)
1613 .unwrap_or_else(|| {
1614 dict_values.push(*pk);
1615 dict_values.len() - 1
1616 });
1617 keys.push(key as u32);
1618 }
1619 let pk_array: ArrayRef = Arc::new(DictionaryArray::<UInt32Type>::new(
1620 UInt32Array::from(keys),
1621 Arc::new(BinaryArray::from_iter_values(dict_values.iter().copied())),
1622 ));
1623
1624 RecordBatch::try_new(
1625 schema,
1626 vec![
1627 Arc::new(UInt64Array::from(field_values.to_vec())),
1628 Arc::new(TimestampMillisecondArray::from_iter_values(
1629 0..primary_keys.len() as i64,
1630 )),
1631 pk_array,
1632 ],
1633 )
1634 .unwrap()
1635 }
1636
1637 fn new_prefilter_batch_binary_pk(primary_keys: &[&[u8]], field_values: &[u64]) -> RecordBatch {
1638 assert_eq!(primary_keys.len(), field_values.len());
1639
1640 let metadata = Arc::new(sst_region_metadata());
1641 let arrow_schema = metadata.schema.arrow_schema();
1642 let field_column = arrow_schema
1643 .field(arrow_schema.index_of("field_0").unwrap())
1644 .clone();
1645 let time_index_column = arrow_schema
1646 .field(arrow_schema.index_of("ts").unwrap())
1647 .clone();
1648 let schema = Arc::new(Schema::new(vec![
1649 field_column,
1650 time_index_column,
1651 Field::new(PRIMARY_KEY_COLUMN_NAME, DataType::Binary, false),
1652 ]));
1653
1654 let pk_array: ArrayRef =
1655 Arc::new(BinaryArray::from_iter_values(primary_keys.iter().copied()));
1656
1657 RecordBatch::try_new(
1658 schema,
1659 vec![
1660 Arc::new(UInt64Array::from(field_values.to_vec())),
1661 Arc::new(TimestampMillisecondArray::from_iter_values(
1662 0..primary_keys.len() as i64,
1663 )),
1664 pk_array,
1665 ],
1666 )
1667 .unwrap()
1668 }
1669
1670 fn field_values(batch: &RecordBatch) -> Vec<u64> {
1671 batch
1672 .column(0)
1673 .as_any()
1674 .downcast_ref::<UInt64Array>()
1675 .unwrap()
1676 .values()
1677 .to_vec()
1678 }
1679
1680 fn remaining_simple_filter_columns(filters: &[SimpleFilterContext]) -> Vec<&str> {
1681 filters
1682 .iter()
1683 .map(|filter_ctx| filter_ctx.filter().as_filter().unwrap().column_name())
1684 .collect()
1685 }
1686
1687 #[test]
1688 fn test_prefilter_primary_key_drops_single_dictionary_batch() {
1689 let metadata = Arc::new(sst_region_metadata());
1690 let filters = Arc::new(new_test_filters(&[col("tag_0").eq(lit("b"))]));
1691 let mut primary_key_filter =
1692 build_primary_key_codec(metadata.as_ref()).primary_key_filter(&metadata, filters);
1693 let pk_a = new_primary_key(&["a", "x"]);
1694 let batch = new_raw_batch(&[pk_a.as_slice(), pk_a.as_slice()], &[10, 11]);
1695 let pk_col_idx = primary_key_column_index(batch.num_columns());
1696
1697 let filtered =
1698 prefilter_flat_batch_by_primary_key(batch, pk_col_idx, primary_key_filter.as_mut())
1699 .unwrap();
1700
1701 assert!(filtered.is_none());
1702 }
1703
1704 #[test]
1705 fn test_prefilter_primary_key_builds_mask_for_fragmented_matches() {
1706 let metadata = Arc::new(sst_region_metadata());
1707 let filters = Arc::new(new_test_filters(&[col("tag_0")
1708 .eq(lit("a"))
1709 .or(col("tag_0").eq(lit("c")))]));
1710 let mut primary_key_filter =
1711 build_primary_key_codec(metadata.as_ref()).primary_key_filter(&metadata, filters);
1712 let pk_a = new_primary_key(&["a", "x"]);
1713 let pk_b = new_primary_key(&["b", "x"]);
1714 let pk_c = new_primary_key(&["c", "x"]);
1715 let pk_d = new_primary_key(&["d", "x"]);
1716 let batch = new_raw_batch(
1717 &[
1718 pk_a.as_slice(),
1719 pk_a.as_slice(),
1720 pk_b.as_slice(),
1721 pk_b.as_slice(),
1722 pk_c.as_slice(),
1723 pk_c.as_slice(),
1724 pk_d.as_slice(),
1725 pk_d.as_slice(),
1726 ],
1727 &[10, 11, 12, 13, 14, 15, 16, 17],
1728 );
1729 let pk_col_idx = primary_key_column_index(batch.num_columns());
1730
1731 let filtered =
1732 prefilter_flat_batch_by_primary_key(batch, pk_col_idx, primary_key_filter.as_mut())
1733 .unwrap()
1734 .unwrap();
1735
1736 assert_eq!(filtered.num_rows(), 4);
1737 assert_eq!(field_values(&filtered), vec![10, 11, 14, 15]);
1738 }
1739
1740 #[test]
1741 fn test_prefilter_builder_returns_none_without_selected_filters() {
1742 let metadata: RegionMetadataRef =
1743 Arc::new(sst_region_metadata_with_encoding(PrimaryKeyEncoding::Dense));
1744 let read_format = FlatReadFormat::new(
1745 metadata.clone(),
1746 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
1747 None,
1748 "test",
1749 false,
1750 )
1751 .unwrap();
1752 let codec = build_primary_key_codec(metadata.as_ref());
1753
1754 let builder = PrefilterContextBuilder::new(
1755 &read_format,
1756 &codec,
1757 None,
1758 None,
1759 Vec::new(),
1760 Vec::new(),
1761 metadata.schema_version,
1762 &stats_metadata(&[vec![1]], None),
1763 );
1764 assert!(builder.is_none());
1765 }
1766
1767 #[test]
1768 fn test_simple_filter_stats_uses_real_metadata() {
1769 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1770 let read_format = FlatReadFormat::new(
1771 metadata.clone(),
1772 ReadColumns::new(
1773 metadata
1774 .column_metadatas
1775 .iter()
1776 .map(|column| column.column_id),
1777 ),
1778 None,
1779 "test",
1780 true,
1781 )
1782 .unwrap();
1783 let parquet_metadata = stats_metadata(&[vec![i64::MAX as u64 + 2]], None);
1784 let filters = new_simple_filter_contexts(
1785 &metadata,
1786 &[
1787 col("field_0").gt(lit(i64::MAX as u64 + 1)),
1788 col("field_0").gt(lit(i64::MAX as u64 + 2)),
1789 lit(i64::MAX as u64 + 1).lt(col("field_0")),
1790 ],
1791 );
1792
1793 let proofs =
1794 simple_filter_stats_proofs(&read_format, parquet_metadata.row_groups(), &filters);
1795 assert_eq!(proofs, vec![vec![true, false, true]]);
1796 }
1797
1798 #[test]
1799 fn test_simple_filter_stats_retain_narrow_or_incomplete_metadata() {
1800 let parquet_metadata = stats_metadata(&[vec![2]], Some(&[true]));
1801 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1802 let read_format = FlatReadFormat::new(
1803 metadata.clone(),
1804 ReadColumns::new(
1805 metadata
1806 .column_metadatas
1807 .iter()
1808 .map(|column| column.column_id),
1809 ),
1810 None,
1811 "test",
1812 true,
1813 )
1814 .unwrap();
1815 let filter = new_simple_filter_contexts(&metadata, &[col("field_0").gt(lit(1_u64))]);
1816 assert_eq!(
1817 simple_filter_stats_proofs(&read_format, parquet_metadata.row_groups(), &filter),
1818 vec![vec![false]],
1819 );
1820 let missing_stats = stats_metadata_with_options(
1821 &[vec![2]],
1822 None,
1823 Some(
1824 parquet::file::properties::WriterProperties::builder()
1825 .set_statistics_enabled(parquet::file::properties::EnabledStatistics::None)
1826 .build(),
1827 ),
1828 );
1829 assert_eq!(
1830 simple_filter_stats_proofs(&read_format, missing_stats.row_groups(), &filter),
1831 vec![vec![false]],
1832 );
1833
1834 let narrow_metadata = metadata_with_field_type(ConcreteDataType::int8_datatype());
1837 let narrow_read_format = FlatReadFormat::new(
1838 narrow_metadata.clone(),
1839 ReadColumns::new(
1840 narrow_metadata
1841 .column_metadatas
1842 .iter()
1843 .map(|column| column.column_id),
1844 ),
1845 None,
1846 "test",
1847 true,
1848 )
1849 .unwrap();
1850 let narrow_filter =
1851 new_simple_filter_contexts(&narrow_metadata, &[col("field_0").gt(lit(1_i32))]);
1852 let narrow_stats = int32_stats_metadata(&[2]);
1853 assert_eq!(
1854 simple_filter_stats_proofs(
1855 &narrow_read_format,
1856 narrow_stats.row_groups(),
1857 &narrow_filter,
1858 ),
1859 vec![vec![false]],
1860 );
1861 }
1862
1863 #[test]
1864 fn test_prefilter_builder_uses_row_group_proofs_by_index() {
1865 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1866 let read_format = FlatReadFormat::new(
1867 metadata.clone(),
1868 ReadColumns::new(
1869 metadata
1870 .column_metadatas
1871 .iter()
1872 .map(|column| column.column_id),
1873 ),
1874 None,
1875 "test",
1876 true,
1877 )
1878 .unwrap();
1879 let codec = build_primary_key_codec(metadata.as_ref());
1880 let builder = PrefilterContextBuilder::new(
1881 &read_format,
1882 &codec,
1883 None,
1884 None,
1885 new_simple_filter_contexts(&metadata, &[col("field_0").gt(lit(1_u64))]),
1886 Vec::new(),
1887 metadata.schema_version,
1888 &stats_metadata(&[vec![2], vec![1]], None),
1889 )
1890 .unwrap();
1891
1892 assert!(builder.build(0).proven_simple_filters[0]);
1893 assert!(!builder.build(1).proven_simple_filters[0]);
1894 }
1895
1896 #[test]
1897 fn test_simple_filter_stats_prove_integer_predicates() {
1898 macro_rules! assert_all_operators {
1899 ($literal:expr, $base:expr, $below:expr, $above:expr) => {{
1900 for (expr, min, max) in [
1901 (col("x").gt(lit($literal)), $above.clone(), $above.clone()),
1902 (col("x").gt_eq(lit($literal)), $base.clone(), $above.clone()),
1903 (col("x").lt(lit($literal)), $below.clone(), $below.clone()),
1904 (col("x").lt_eq(lit($literal)), $below.clone(), $base.clone()),
1905 (col("x").eq(lit($literal)), $base.clone(), $base.clone()),
1906 (
1907 col("x").not_eq(lit($literal)),
1908 $below.clone(),
1909 $below.clone(),
1910 ),
1911 ] {
1912 let filter = SimpleFilterEvaluator::try_new(&expr).unwrap();
1913 assert!(simple_filter_is_true_by_values(
1914 &filter,
1915 &$base,
1916 Some(min),
1917 Some(max),
1918 Some(Value::UInt64(0)),
1919 ));
1920 }
1921 }};
1922 }
1923
1924 assert_all_operators!(5_i32, Value::Int32(5), Value::Int32(-1), Value::Int32(6));
1925 assert_all_operators!(5_u32, Value::UInt32(5), Value::UInt32(4), Value::UInt32(6));
1926 assert_all_operators!(5_i64, Value::Int64(5), Value::Int64(-1), Value::Int64(6));
1927 assert_all_operators!(
1928 i64::MAX as u64 + 2,
1929 Value::UInt64(i64::MAX as u64 + 2),
1930 Value::UInt64(i64::MAX as u64 + 1),
1931 Value::UInt64(i64::MAX as u64 + 3)
1932 );
1933
1934 for (expr, min, max) in [
1937 (col("x").lt(lit(5_i64)), 1, 5),
1938 (col("x").gt(lit(5_i64)), 5, 9),
1939 (col("x").lt_eq(lit(5_i64)), 1, 6),
1940 (col("x").gt_eq(lit(5_i64)), 4, 9),
1941 (col("x").eq(lit(5_i64)), 5, 6),
1942 (col("x").not_eq(lit(5_i64)), 5, 9),
1943 (col("x").not_eq(lit(5_i64)), 1, 5),
1944 (col("x").not_eq(lit(5_i64)), 1, 9),
1945 ] {
1946 let filter = SimpleFilterEvaluator::try_new(&expr).unwrap();
1947 assert!(
1948 !simple_filter_is_true_by_values(
1949 &filter,
1950 &Value::Int64(5),
1951 Some(Value::Int64(min)),
1952 Some(Value::Int64(max)),
1953 Some(Value::UInt64(0)),
1954 ),
1955 "{expr:?} must not be proven by stats {min}..={max}",
1956 );
1957 }
1958 }
1959
1960 #[test]
1961 fn test_simple_filter_stats_retain_unknown_or_unsupported_values() {
1962 let filter = SimpleFilterEvaluator::try_new(&col("x").gt(lit(1_i32))).unwrap();
1963 let literal = Value::Int32(1);
1964 for (min, max, null_count) in [
1965 (
1966 Some(Value::Int32(2)),
1967 Some(Value::Int32(3)),
1968 Some(Value::UInt64(1)),
1969 ),
1970 (Some(Value::Int32(2)), Some(Value::Int32(3)), None),
1971 (
1972 Some(Value::Null),
1973 Some(Value::Int32(3)),
1974 Some(Value::UInt64(0)),
1975 ),
1976 (None, Some(Value::Int32(3)), Some(Value::UInt64(0))),
1977 (Some(Value::Int32(2)), None, Some(Value::UInt64(0))),
1978 (
1979 Some(Value::Int64(2)),
1980 Some(Value::Int64(3)),
1981 Some(Value::UInt64(0)),
1982 ),
1983 (
1984 Some(Value::Int32(3)),
1985 Some(Value::Int32(2)),
1986 Some(Value::UInt64(0)),
1987 ),
1988 ] {
1989 assert!(!simple_filter_is_true_by_values(
1990 &filter, &literal, min, max, null_count
1991 ));
1992 }
1993
1994 let float_filter = SimpleFilterEvaluator::try_new(&col("x").gt(lit(1.0_f64))).unwrap();
1995 assert!(!simple_filter_is_true_by_values(
1996 &float_filter,
1997 &Value::Float64(1.0.into()),
1998 Some(Value::Float64(2.0.into())),
1999 Some(Value::Float64(3.0.into())),
2000 Some(Value::UInt64(0)),
2001 ));
2002 }
2003
2004 #[test]
2005 fn test_prefilter_entries_keep_only_unproven_simple_filter_indices() {
2006 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2007 let filters = new_simple_filter_contexts(
2008 &metadata,
2009 &[col("field_0").gt(lit(1_u64)), col("field_0").lt(lit(9_u64))],
2010 );
2011 let context = PrefilterContext {
2012 pk_filter: None,
2013 filters,
2014 physical_filters: Vec::new(),
2015 schema_version: metadata.schema_version,
2016 pk_filter_expr_strs: None,
2017 arrow_schema: metadata.schema.arrow_schema().clone(),
2018 proven_simple_filters: vec![true, false],
2019 };
2020
2021 let entries = all_prefilter_entries(&context);
2022 assert!(matches!(
2023 entries.as_slice(),
2024 [PrefilterEntry {
2025 kind: PrefilterEntryKind::Simple(1),
2026 ..
2027 }]
2028 ));
2029 }
2030
2031 #[test]
2032 fn test_identity_row_selection_preserves_input() {
2033 let sparse = RowSelection::from(vec![
2034 RowSelector::skip(2),
2035 RowSelector::select(3),
2036 RowSelector::skip(1),
2037 ]);
2038 assert_eq!(identity_row_selection(&Some(sparse.clone()), 6), sparse);
2039
2040 let empty = RowSelection::from(vec![]);
2041 assert_eq!(identity_row_selection(&Some(empty.clone()), 6), empty);
2042 assert_eq!(
2043 identity_row_selection(&None, 6),
2044 RowSelection::from(vec![RowSelector::select(6)])
2045 );
2046 }
2047
2048 #[test]
2049 fn test_should_use_prefilter() {
2050 assert!(should_use_prefilter(1, 5, 6));
2051 assert!(!should_use_prefilter(1, 0, 1));
2052 assert!(!should_use_prefilter(1, 1, 2));
2053 assert!(!should_use_prefilter(4, 3, 7));
2054 assert!(should_use_prefilter(3, 3, 6));
2055 }
2056
2057 #[test]
2058 fn test_build_bulk_filter_plan_classifies_filters_across_read_paths() {
2059 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata_with_encoding(
2060 PrimaryKeyEncoding::Sparse,
2061 ));
2062 let legacy_read_format = FlatReadFormat::new(
2063 metadata.clone(),
2064 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
2065 None,
2066 "memtable",
2067 false,
2068 )
2069 .unwrap();
2070 assert!(!legacy_read_format.batch_has_raw_pk_columns());
2071
2072 let plan = build_bulk_filter_plan(
2073 &legacy_read_format,
2074 Some(&Predicate::new(vec![
2075 col("tag_0").eq(lit("a")),
2076 col("field_0").gt(lit(1_u64)),
2077 ])),
2078 );
2079 assert_eq!(
2080 plan.pk_filters.as_ref().map(|filters| filters.len()),
2081 Some(1)
2082 );
2083 assert_eq!(
2084 remaining_simple_filter_columns(&plan.remaining_simple_filters),
2085 vec!["field_0"]
2086 );
2087
2088 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2089 let raw_pk_read_format = FlatReadFormat::new(
2090 metadata.clone(),
2091 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
2092 None,
2093 "memtable",
2094 true,
2095 )
2096 .unwrap();
2097 assert!(raw_pk_read_format.batch_has_raw_pk_columns());
2098
2099 let tag_only_plan = build_bulk_filter_plan(
2100 &raw_pk_read_format,
2101 Some(&Predicate::new(vec![col("tag_0").eq(lit("a"))])),
2102 );
2103 assert!(tag_only_plan.pk_filters.is_none());
2104 assert_eq!(
2105 remaining_simple_filter_columns(&tag_only_plan.remaining_simple_filters),
2106 vec!["tag_0"]
2107 );
2108
2109 let field_only_plan = build_bulk_filter_plan(
2110 &raw_pk_read_format,
2111 Some(&Predicate::new(vec![col("field_0").gt(lit(1_u64))])),
2112 );
2113 assert!(field_only_plan.pk_filters.is_none());
2114 assert_eq!(
2115 remaining_simple_filter_columns(&field_only_plan.remaining_simple_filters),
2116 vec!["field_0"]
2117 );
2118 }
2119
2120 #[test]
2121 fn test_build_reader_filter_plan_classifies_filters_for_prefilter_modes() {
2122 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2123 let full_read_format = FlatReadFormat::new(
2124 metadata.clone(),
2125 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
2126 None,
2127 "test",
2128 true,
2129 )
2130 .unwrap();
2131 let codec = build_primary_key_codec(metadata.as_ref());
2132
2133 let skip_fields_plan = build_reader_filter_plan(
2134 Some(&Predicate::new(vec![
2135 col("tag_0").eq(lit("a")),
2136 col("field_0").gt(lit(1_u64)),
2137 ])),
2138 None,
2139 PreFilterMode::SkipFields,
2140 true,
2141 false,
2142 &full_read_format,
2143 &codec,
2144 &stats_metadata(&[vec![1]], None),
2145 );
2146 assert!(skip_fields_plan.prefilter_builder.is_some());
2147 assert_eq!(
2148 remaining_simple_filter_columns(&skip_fields_plan.remaining_simple_filters),
2149 vec!["field_0"]
2150 );
2151
2152 let postponed_time_plan = build_reader_filter_plan(
2153 Some(&Predicate::new(vec![
2154 col("tag_0").eq(lit("a")),
2155 col("field_0").gt(lit(1_u64)),
2156 col("ts").gt_eq(lit(ScalarValue::TimestampMillisecond(Some(1), None))),
2157 ])),
2158 None,
2159 PreFilterMode::SkipFields,
2160 true,
2161 true,
2162 &full_read_format,
2163 &codec,
2164 &stats_metadata(&[vec![1]], None),
2165 );
2166 assert!(postponed_time_plan.prefilter_builder.is_some());
2167 assert_eq!(
2168 remaining_simple_filter_columns(&postponed_time_plan.remaining_simple_filters),
2169 vec!["field_0", "ts"]
2170 );
2171
2172 let postponed_time_only_plan = build_reader_filter_plan(
2173 Some(&Predicate::new(vec![col("ts").gt_eq(lit(
2174 ScalarValue::TimestampMillisecond(Some(1), None),
2175 ))])),
2176 None,
2177 PreFilterMode::All,
2178 true,
2179 true,
2180 &full_read_format,
2181 &codec,
2182 &stats_metadata(&[vec![1]], None),
2183 );
2184 assert!(postponed_time_only_plan.prefilter_builder.is_none());
2185 assert_eq!(
2186 remaining_simple_filter_columns(&postponed_time_only_plan.remaining_simple_filters),
2187 vec!["ts"]
2188 );
2189
2190 let metric_metadata: RegionMetadataRef = Arc::new(sst_region_metadata_with_encoding(
2191 PrimaryKeyEncoding::Sparse,
2192 ));
2193 let field_0 = metric_metadata.column_by_name("field_0").unwrap().column_id;
2194 let ts = metric_metadata.time_index_column().column_id;
2195 let projected_read_format = FlatReadFormat::new(
2196 metric_metadata.clone(),
2197 ReadColumns::new([field_0, ts]),
2198 None,
2199 "test",
2200 true,
2201 )
2202 .unwrap();
2203 let metric_codec = build_primary_key_codec(metric_metadata.as_ref());
2204 let pk_prefilter_plan = build_reader_filter_plan(
2205 Some(&Predicate::new(vec![col("tag_0").eq(lit("a"))])),
2206 None,
2207 PreFilterMode::All,
2208 true,
2209 false,
2210 &projected_read_format,
2211 &metric_codec,
2212 &stats_metadata(&[vec![1]], None),
2213 );
2214 assert!(pk_prefilter_plan.prefilter_builder.is_some());
2215 assert!(
2216 pk_prefilter_plan
2217 .prefilter_builder
2218 .as_ref()
2219 .unwrap()
2220 .build_primary_key_filter()
2221 .is_some()
2222 );
2223 assert!(pk_prefilter_plan.remaining_simple_filters.is_empty());
2224
2225 let disabled_plan = build_reader_filter_plan(
2226 Some(&Predicate::new(vec![
2227 col("tag_0").eq(lit("a")),
2228 col("field_0").gt(lit(1_u64)),
2229 col("ts").gt_eq(lit(ScalarValue::TimestampMillisecond(Some(1), None))),
2230 ])),
2231 None,
2232 PreFilterMode::All,
2233 false,
2234 true,
2235 &projected_read_format,
2236 &metric_codec,
2237 &stats_metadata(&[vec![1]], None),
2238 );
2239 assert!(disabled_plan.prefilter_builder.is_none());
2240 assert_eq!(
2241 remaining_simple_filter_columns(&disabled_plan.remaining_simple_filters),
2242 vec!["tag_0", "field_0", "ts"]
2243 );
2244 }
2245
2246 #[test]
2247 fn test_pk_filter_expr_strings_are_stable_under_expr_order() {
2248 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata_with_encoding(
2249 PrimaryKeyEncoding::Sparse,
2250 ));
2251 let read_format = FlatReadFormat::new(
2252 metadata.clone(),
2253 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
2254 None,
2255 "test",
2256 false,
2257 )
2258 .unwrap();
2259 let codec = build_primary_key_codec(metadata.as_ref());
2260
2261 let expr_a = col("tag_0").eq(lit("a"));
2262 let expr_b = col("tag_1").eq(lit("x"));
2263 let plan_ab = build_reader_filter_plan(
2264 Some(&Predicate::new(vec![expr_a.clone(), expr_b.clone()])),
2265 None,
2266 PreFilterMode::All,
2267 true,
2268 false,
2269 &read_format,
2270 &codec,
2271 &stats_metadata(&[vec![1]], None),
2272 );
2273 let plan_b_a = build_reader_filter_plan(
2274 Some(&Predicate::new(vec![expr_b, expr_a])),
2275 None,
2276 PreFilterMode::All,
2277 true,
2278 false,
2279 &read_format,
2280 &codec,
2281 &stats_metadata(&[vec![1]], None),
2282 );
2283
2284 let exprs_ab = plan_ab.prefilter_builder.unwrap().pk_filter_expr_strs;
2285 let exprs_b_a = plan_b_a.prefilter_builder.unwrap().pk_filter_expr_strs;
2286 assert!(exprs_ab.is_some());
2287 assert_eq!(exprs_ab, exprs_b_a);
2288 }
2289
2290 #[test]
2291 fn test_simple_and_physical_contexts_preserve_expr_strings() {
2292 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2293 let read_format = FlatReadFormat::new(
2294 metadata.clone(),
2295 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
2296 None,
2297 "test",
2298 true,
2299 )
2300 .unwrap();
2301
2302 let simple_expr = col("tag_0").eq(lit("a"));
2303 let simple = SimpleFilterContext::new_opt(&metadata, None, &simple_expr).unwrap();
2304 assert_eq!(simple.expr_str(), format!("{simple_expr:?}"));
2305
2306 let physical_expr = col("field_0").in_list(vec![lit(1_u64), lit(2_u64)], false);
2307 let physical =
2308 PhysicalFilterContext::new_opt(&metadata, None, &read_format, &physical_expr).unwrap();
2309 assert_eq!(physical.expr_str(), format!("{physical_expr:?}"));
2310 }
2311
2312 #[test]
2313 fn test_eval_simple_filter_mask_uses_flat_tag_columns_directly() {
2314 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2315 let filters = new_simple_filter_contexts(&metadata, &[col("tag_0").eq(lit("a"))]);
2316 let batch = new_record_batch_with_custom_sequence(&["a", "x"], 0, 4, 1);
2317
2318 let mask = eval_simple_filter_mask(&batch, &filters[0], "test").unwrap();
2319 assert_eq!(mask.count_set_bits(), 4);
2320 }
2321
2322 #[test]
2323 fn test_eval_simple_filter_mask_errors_on_missing_selected_column() {
2324 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2325 let filters = new_simple_filter_contexts(&metadata, &[col("tag_0").eq(lit("a"))]);
2326 let pk = new_primary_key(&["a", "x"]);
2327 let batch = new_raw_batch(&[pk.as_slice()], &[10]);
2328
2329 let err = eval_simple_filter_mask(&batch, &filters[0], "test").unwrap_err();
2330 let err = err.to_string();
2331 assert!(err.contains("Prefilter column"));
2332 assert!(err.contains("tag_0"));
2333 }
2334
2335 #[test]
2336 fn test_eval_physical_filter_mask_evaluates_physical_filters() {
2337 let metadata: RegionMetadataRef =
2338 Arc::new(sst_region_metadata_with_encoding(PrimaryKeyEncoding::Dense));
2339 let read_format = FlatReadFormat::new(
2340 metadata.clone(),
2341 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
2342 None,
2343 "test",
2344 false,
2345 )
2346 .unwrap();
2347 let expr = col("field_0").in_list(vec![lit(11_u64)], false);
2348 let physical_filters = new_physical_filter_contexts(&metadata, &read_format, &[expr]);
2349 let pk = new_primary_key(&["a", "x"]);
2350 let batch = new_raw_batch(&[pk.as_slice(), pk.as_slice(), pk.as_slice()], &[9, 10, 11]);
2351
2352 let mask = eval_physical_filter_mask(&batch, &physical_filters[0], "test").unwrap();
2353 assert_eq!(mask.count_set_bits(), 1);
2354 }
2355
2356 #[test]
2357 fn test_eval_pk_group_mask_finds_pk_column_by_name() {
2358 let metadata = Arc::new(sst_region_metadata());
2359 let filters = Arc::new(new_test_filters(&[col("tag_0").eq(lit("a"))]));
2360 let mut pk_filter = Some(Box::new(CachedPrimaryKeyFilter::new(
2361 build_primary_key_codec(metadata.as_ref()).primary_key_filter(&metadata, filters),
2362 )) as Box<dyn PrimaryKeyFilter>);
2363 let pk_a = new_primary_key(&["a", "x"]);
2364 let pk_b = new_primary_key(&["b", "x"]);
2365 let batch = new_prefilter_batch(
2366 &[
2367 pk_a.as_slice(),
2368 pk_a.as_slice(),
2369 pk_b.as_slice(),
2370 pk_b.as_slice(),
2371 ],
2372 &[10, 11, 12, 13],
2373 );
2374
2375 let mask = primary_key_filter_mask(&batch, pk_filter.as_mut().unwrap().as_mut()).unwrap();
2376
2377 assert_eq!(mask.count_set_bits(), 2);
2378 }
2379
2380 #[test]
2381 fn test_eval_pk_group_mask_handles_binary_pk_column() {
2382 let metadata = Arc::new(sst_region_metadata());
2383 let filters = Arc::new(new_test_filters(&[col("tag_0").eq(lit("a"))]));
2384 let mut pk_filter = Some(Box::new(CachedPrimaryKeyFilter::new(
2385 build_primary_key_codec(metadata.as_ref()).primary_key_filter(&metadata, filters),
2386 )) as Box<dyn PrimaryKeyFilter>);
2387 let pk_a = new_primary_key(&["a", "x"]);
2388 let pk_b = new_primary_key(&["b", "x"]);
2389 let batch = new_prefilter_batch_binary_pk(
2390 &[
2391 pk_a.as_slice(),
2392 pk_a.as_slice(),
2393 pk_b.as_slice(),
2394 pk_b.as_slice(),
2395 ],
2396 &[10, 11, 12, 13],
2397 );
2398
2399 let mask = primary_key_filter_mask(&batch, pk_filter.as_mut().unwrap().as_mut()).unwrap();
2400
2401 assert_eq!(mask.count_set_bits(), 2);
2402 }
2403}