Skip to main content

mito2/sst/parquet/
file_range.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Structs and functions for reading ranges from a parquet file. A file range
16//! is usually a row group in a parquet file.
17
18use std::collections::{HashMap, HashSet};
19use std::ops::{BitAnd, Range};
20use std::sync::Arc;
21
22use api::v1::{OpType, SemanticType};
23use common_telemetry::error;
24use datafusion::physical_plan::PhysicalExpr;
25use datafusion::physical_plan::expressions::DynamicFilterPhysicalExpr;
26use datatypes::arrow::array::{Array as _, ArrayRef, BooleanArray};
27use datatypes::arrow::buffer::BooleanBuffer;
28use datatypes::arrow::record_batch::RecordBatch;
29use datatypes::schema::Schema;
30use futures::StreamExt;
31use mito_codec::row_converter::PrimaryKeyCodec;
32use object_store::ObjectStore;
33use parquet::arrow::arrow_reader::RowSelection;
34use parquet::file::metadata::ParquetMetaData;
35use parquet::file::statistics::Statistics;
36use snafu::{OptionExt, ResultExt, ensure};
37use store_api::codec::PrimaryKeyEncoding;
38use store_api::metadata::RegionMetadataRef;
39use store_api::storage::{ColumnId, TimeSeriesRowSelector};
40use table::predicate::Predicate;
41use tokio::sync::OnceCell;
42
43use crate::cache::CacheStrategy;
44use crate::error::{
45    ComputeArrowSnafu, DecodeStatsSnafu, EvalPartitionFilterSnafu, InvalidRecordBatchSnafu,
46    NewRecordBatchSnafu, RecordBatchSnafu, Result, StatsNotPresentSnafu, UnexpectedSnafu,
47};
48use crate::read::compat::FlatCompatBatch;
49use crate::read::flat_projection::CompactionProjectionMapper;
50use crate::read::last_row::FlatRowGroupLastRowCachedReader;
51use crate::read::prune::FlatPruneReader;
52use crate::sst::file::FileHandle;
53use crate::sst::parquet::flat_format::{
54    DecodedPrimaryKeys, FlatReadFormat, decode_primary_keys, primary_key_column_index,
55    time_index_column_index,
56};
57use crate::sst::parquet::json_align::ProjectedRecordBatchStream;
58use crate::sst::parquet::prefilter::primary_key_filter_mask;
59use crate::sst::parquet::reader::{
60    FlatRowGroupReader, MaybeFilter, RowGroupBuildContext, RowGroupReaderBuilder,
61    SimpleFilterContext,
62};
63use crate::sst::parquet::row_group::ParquetFetchMetrics;
64use crate::sst::parquet::row_selection::{intersect_row_selections, row_selection_from_row_ranges};
65use crate::sst::parquet::stats::RowGroupPruningStats;
66use crate::sst::range_index::{SstRangeIndexSearcher, range_index_path};
67
68/// Checks if a row group contains delete operations by examining the min value of op_type column.
69///
70/// Returns `Ok(true)` if the row group contains delete operations, `Ok(false)` if it doesn't,
71/// or an error if the statistics are not present or cannot be decoded.
72pub(crate) fn row_group_contains_delete(
73    parquet_meta: &ParquetMetaData,
74    row_group_index: usize,
75    file_path: &str,
76) -> Result<bool> {
77    let row_group_metadata = &parquet_meta.row_groups()[row_group_index];
78
79    // safety: The last column of SST must be op_type
80    let column_metadata = &row_group_metadata.columns().last().unwrap();
81    let stats = column_metadata
82        .statistics()
83        .context(StatsNotPresentSnafu { file_path })?;
84    stats
85        .min_bytes_opt()
86        .context(StatsNotPresentSnafu { file_path })?
87        .try_into()
88        .map(i32::from_le_bytes)
89        .map(|min_op_type| min_op_type == OpType::Delete as i32)
90        .ok()
91        .context(DecodeStatsSnafu { file_path })
92}
93
94/// A range of a parquet SST. Now it is a row group.
95/// We can read different file ranges in parallel.
96#[derive(Clone)]
97pub struct FileRange {
98    /// Shared context.
99    context: FileRangeContextRef,
100    /// Index of the row group in the SST.
101    row_group_idx: usize,
102    /// Row selection for the row group. `None` means all rows.
103    row_selection: Option<RowSelection>,
104}
105
106impl FileRange {
107    /// Returns the shared range-index searcher, opening it on first use.
108    pub(crate) async fn range_index_searcher(&self) -> Result<Option<&SstRangeIndexSearcher>> {
109        self.context.range_index_searcher().await
110    }
111
112    /// Returns the region metadata stored in this SST.
113    pub(crate) fn region_metadata(&self) -> &RegionMetadataRef {
114        self.context.read_format().metadata()
115    }
116
117    /// Returns encoded primary-key min/max statistics for this row group.
118    pub(crate) fn primary_key_range(&self) -> Option<(&[u8], &[u8])> {
119        let metadata = self.context.reader_builder.parquet_metadata();
120        let num_columns = metadata.file_metadata().schema_descr().num_columns();
121        let primary_key_index = primary_key_column_index(num_columns);
122        match metadata
123            .row_group(self.row_group_idx)
124            .column(primary_key_index)
125            .statistics()?
126        {
127            Statistics::ByteArray(statistics) => {
128                Some((statistics.min_bytes_opt()?, statistics.max_bytes_opt()?))
129            }
130            _ => None,
131        }
132    }
133
134    /// Creates a new [FileRange].
135    pub(crate) fn new(
136        context: FileRangeContextRef,
137        row_group_idx: usize,
138        row_selection: Option<RowSelection>,
139    ) -> Self {
140        Self {
141            context,
142            row_group_idx,
143            row_selection,
144        }
145    }
146
147    /// Returns true if [FileRange] selects all rows in row group.
148    fn select_all(&self) -> bool {
149        let rows_in_group = self
150            .context
151            .reader_builder
152            .parquet_metadata()
153            .row_group(self.row_group_idx)
154            .num_rows();
155
156        let Some(row_selection) = &self.row_selection else {
157            return true;
158        };
159        row_selection.row_count() == rows_in_group as usize
160    }
161
162    /// Performs pruning before reading the [FileRange].
163    /// It use latest dynamic filters with row group statistics to prune the range.
164    ///
165    /// Returns false if the entire range is pruned and can be skipped.
166    fn in_dynamic_filter_range(&self) -> bool {
167        if self.context.base.dyn_filters.is_empty() {
168            return true;
169        }
170        let curr_row_group = self
171            .context
172            .reader_builder
173            .parquet_metadata()
174            .row_group(self.row_group_idx);
175        let read_format = self.context.read_format();
176        let prune_schema = &self.context.base.prune_schema;
177        let stats = RowGroupPruningStats::new(
178            std::slice::from_ref(curr_row_group),
179            read_format,
180            self.context.base.expected_metadata.clone(),
181            self.context.base.pre_filter_mode.skip_fields(),
182        );
183
184        // not costly to create a predicate here since dynamic filters are wrapped in Arc
185        let pred = Predicate::with_dyn_filters(vec![], self.context.base.dyn_filters.clone());
186
187        pred.prune_with_stats(&stats, prune_schema.arrow_schema())
188            .first()
189            .cloned()
190            .unwrap_or(true) // unexpected, not skip just in case
191    }
192
193    /// Creates a flat reader that returns RecordBatch.
194    pub async fn flat_reader(
195        &self,
196        selector: Option<TimeSeriesRowSelector>,
197        fetch_metrics: Option<&ParquetFetchMetrics>,
198    ) -> Result<Option<FlatPruneReader>> {
199        if !self.in_dynamic_filter_range() {
200            return Ok(None);
201        }
202        // Compute skip_fields once for this row group
203        let skip_fields = self.context.base.pre_filter_mode.skip_fields();
204        let parquet_reader = self
205            .context
206            .reader_builder
207            .build(self.context.build_context(
208                self.row_group_idx,
209                self.row_selection.clone(),
210                fetch_metrics,
211            ))
212            .await?;
213
214        let use_last_row_reader = if matches!(
215            selector,
216            Some(TimeSeriesRowSelector::LastRow { after_merge: false })
217        ) {
218            // Only use the last row reader if row group does not contain DELETE, all
219            // rows are selected, and filters that still run after this reader
220            // cannot change which row is last. Tag filters are safe because a
221            // tag is constant within a series. Timestamp and field filters are
222            // not safe for this shortcut.
223            let put_only = !self
224                .context
225                .contains_delete(self.row_group_idx)
226                .inspect_err(|e| {
227                    error!(e; "Failed to decode min value of op_type, fallback to FlatRowGroupReader");
228                })
229                .unwrap_or(true);
230            put_only && self.select_all() && self.context.remaining_filters_preserve_last_row()
231        } else {
232            false
233        };
234
235        let flat_prune_reader = if use_last_row_reader {
236            let flat_row_group_reader =
237                FlatRowGroupReader::new(self.context.clone(), parquet_reader);
238            // Predicate prefiltering makes the input stream predicate-dependent, so cached
239            // selector results are not reusable across queries with different filters.
240            let cache_strategy = if self.context.reader_builder.has_predicate_prefilter() {
241                CacheStrategy::Disabled
242            } else {
243                self.context.reader_builder.cache_strategy().clone()
244            };
245            let reader = FlatRowGroupLastRowCachedReader::new(
246                self.file_handle().file_id().file_id(),
247                self.row_group_idx,
248                cache_strategy,
249                self.context.read_format().parquet_read_columns(),
250                self.context.read_format().json_target_types().clone(),
251                flat_row_group_reader,
252            );
253            FlatPruneReader::new_with_last_row_reader(self.context.clone(), reader, skip_fields)
254        } else {
255            let flat_row_group_reader =
256                FlatRowGroupReader::new(self.context.clone(), parquet_reader);
257            FlatPruneReader::new_with_row_group_reader(
258                self.context.clone(),
259                flat_row_group_reader,
260                skip_fields,
261            )
262        };
263
264        Ok(Some(flat_prune_reader))
265    }
266
267    /// Creates a reader that returns only the encoded primary-key column.
268    ///
269    /// The returned primary keys are compatible with the expected region metadata.
270    pub(crate) async fn primary_key_reader(
271        &self,
272        fetch_metrics: Option<&ParquetFetchMetrics>,
273    ) -> Result<Option<ProjectedRecordBatchStream>> {
274        self.primary_key_reader_inner(fetch_metrics, true).await
275    }
276
277    async fn primary_key_reader_inner(
278        &self,
279        fetch_metrics: Option<&ParquetFetchMetrics>,
280        check_dynamic_filter: bool,
281    ) -> Result<Option<ProjectedRecordBatchStream>> {
282        if check_dynamic_filter && !self.in_dynamic_filter_range() {
283            return Ok(None);
284        }
285        let stream = self
286            .context
287            .reader_builder
288            .build_primary_key(self.context.build_context(
289                self.row_group_idx,
290                self.row_selection.clone(),
291                fetch_metrics,
292            ))
293            .await?;
294        if self.context.compat_batch().is_none() {
295            return Ok(Some(stream));
296        }
297
298        let context = self.context.clone();
299        let stream = stream
300            .map(move |batch| {
301                let batch = batch?;
302                let compat = context.compat_batch().context(UnexpectedSnafu {
303                    reason: "Primary-key compatibility helper is missing",
304                })?;
305                let primary_key = compat.compat_primary_key(batch.column(0))?;
306                RecordBatch::try_new(batch.schema(), vec![primary_key]).context(NewRecordBatchSnafu)
307            })
308            .boxed();
309        Ok(Some(stream))
310    }
311
312    /// Builds a full-projection reader selected only by the provided encoded-PK
313    /// filter and this range's existing row selection.
314    ///
315    /// This deliberately bypasses generic predicate prefiltering. The series
316    /// pruner selected the row group independently, and simple predicates retained
317    /// by the disabled prefilter plan are applied precisely before merge.
318    pub(crate) async fn reader_by_primary_key(
319        &self,
320        primary_key_filter: &mut dyn mito_codec::row_converter::PrimaryKeyFilter,
321        fetch_metrics: Option<&ParquetFetchMetrics>,
322    ) -> Result<Option<FlatRowGroupReader>> {
323        let Some(mut primary_keys) = self.primary_key_reader_inner(fetch_metrics, false).await?
324        else {
325            return Ok(None);
326        };
327
328        let mut masks = Vec::new();
329        while let Some(batch) = primary_keys.next().await {
330            let batch = batch?;
331            masks.push(BooleanArray::from(primary_key_filter_mask(
332                &batch,
333                primary_key_filter,
334            )?));
335        }
336        let Some(selected) = refine_primary_key_selection(&masks, &self.row_selection) else {
337            return Ok(None);
338        };
339
340        let stream = self
341            .context
342            .reader_builder
343            .build_without_prefilter(self.context.build_context(
344                self.row_group_idx,
345                Some(selected),
346                fetch_metrics,
347            ))
348            .await?;
349        Ok(Some(FlatRowGroupReader::new(self.context.clone(), stream)))
350    }
351
352    /// Returns the source SST row-group index.
353    pub(crate) fn row_group_index(&self) -> usize {
354        self.row_group_idx
355    }
356
357    /// Builds a reader from absolute row-group offsets supplied by a range index.
358    /// Existing pruning is intersected before reading, without predicate prefiltering.
359    pub(crate) async fn reader_by_row_ranges(
360        &self,
361        ranges: Vec<Range<usize>>,
362        fetch_metrics: Option<&ParquetFetchMetrics>,
363    ) -> Result<Option<FlatRowGroupReader>> {
364        let num_rows = self
365            .context
366            .reader_builder
367            .parquet_metadata()
368            .row_group(self.row_group_idx)
369            .num_rows() as usize;
370        let Some(selected) = refine_row_range_selection(ranges, num_rows, &self.row_selection)?
371        else {
372            return Ok(None);
373        };
374        let stream = self
375            .context
376            .reader_builder
377            .build_without_prefilter(self.context.build_context(
378                self.row_group_idx,
379                Some(selected),
380                fetch_metrics,
381            ))
382            .await?;
383        Ok(Some(FlatRowGroupReader::new(self.context.clone(), stream)))
384    }
385
386    /// Returns the helper to compat batches.
387    pub(crate) fn compat_batch(&self) -> Option<&FlatCompatBatch> {
388        self.context.compat_batch()
389    }
390
391    /// Returns the helper to project batches.
392    pub(crate) fn compaction_projection_mapper(&self) -> Option<&CompactionProjectionMapper> {
393        self.context.compaction_projection_mapper()
394    }
395
396    /// Filters a full-projection batch using this range's precise filters.
397    pub(crate) fn precise_filter_flat(
398        &self,
399        input: RecordBatch,
400        skip_fields: bool,
401        skip_tags: bool,
402    ) -> Result<Option<RecordBatch>> {
403        self.context
404            .precise_filter_flat(input, skip_fields, skip_tags)
405    }
406
407    /// Returns the precise-filter mode configured for this range.
408    pub(crate) fn pre_filter_mode(&self) -> PreFilterMode {
409        self.context.pre_filter_mode()
410    }
411
412    /// Returns the file handle of the file range.
413    pub(crate) fn file_handle(&self) -> &FileHandle {
414        self.context.reader_builder.file_handle()
415    }
416}
417
418/// Intersects source-row offsets, rather than offsets within an already selected stream.
419fn refine_row_range_selection(
420    ranges: Vec<Range<usize>>,
421    num_rows: usize,
422    original: &Option<RowSelection>,
423) -> Result<Option<RowSelection>> {
424    let mut previous_end = 0;
425    for range in &ranges {
426        ensure!(
427            range.start >= previous_end && range.start < range.end && range.end <= num_rows,
428            InvalidRecordBatchSnafu {
429                reason: format!(
430                    "invalid range-index row range {range:?} after {previous_end}, row group has {num_rows} rows"
431                ),
432            }
433        );
434        previous_end = range.end;
435    }
436    let selected = row_selection_from_row_ranges(ranges.into_iter(), num_rows);
437    let selected = match original {
438        Some(original) => intersect_row_selections(original, &selected),
439        None => selected,
440    };
441    Ok((selected.row_count() > 0).then_some(selected))
442}
443
444fn refine_primary_key_selection(
445    masks: &[BooleanArray],
446    original: &Option<RowSelection>,
447) -> Option<RowSelection> {
448    if masks.is_empty() {
449        return None;
450    }
451    let selected = RowSelection::from_filters(masks);
452    let selected = match original {
453        Some(original) => original.and_then(&selected),
454        None => selected,
455    };
456    (selected.row_count() > 0).then_some(selected)
457}
458
459/// Context shared by ranges of the same parquet SST.
460pub struct FileRangeContext {
461    /// Store for a range index registered in the scan's index snapshot.
462    range_index_store: Option<ObjectStore>,
463    /// Lazily opened range index shared by all ranges of this file.
464    range_index_searcher: OnceCell<SstRangeIndexSearcher>,
465    /// Row group reader builder for the file.
466    reader_builder: RowGroupReaderBuilder,
467    /// Base of the context.
468    base: RangeBase,
469}
470
471pub type FileRangeContextRef = Arc<FileRangeContext>;
472
473impl FileRangeContext {
474    /// Creates a new [FileRangeContext].
475    pub(crate) fn new(
476        reader_builder: RowGroupReaderBuilder,
477        base: RangeBase,
478        range_index_store: Option<ObjectStore>,
479    ) -> Self {
480        Self {
481            reader_builder,
482            base,
483            range_index_store,
484            range_index_searcher: OnceCell::new(),
485        }
486    }
487
488    /// Opens the range index once, retaining the SST handle throughout its use.
489    async fn range_index_searcher(&self) -> Result<Option<&SstRangeIndexSearcher>> {
490        let Some(store) = &self.range_index_store else {
491            return Ok(None);
492        };
493        self.range_index_searcher
494            .get_or_try_init(|| async {
495                let file = self.reader_builder.file_handle();
496                let path = range_index_path(file.region_id(), file.file_id().file_id());
497                SstRangeIndexSearcher::open(store.clone(), &path).await
498            })
499            .await
500            .map(Some)
501    }
502
503    /// Returns filters pushed down.
504    pub(crate) fn filters(&self) -> &[SimpleFilterContext] {
505        &self.base.filters
506    }
507
508    /// Returns true if a partition filter is configured.
509    pub(crate) fn has_partition_filter(&self) -> bool {
510        self.base.partition_filter.is_some()
511    }
512
513    /// Returns true if applying the remaining precise filters after selecting
514    /// the last row cannot change which row is selected for a series.
515    fn remaining_filters_preserve_last_row(&self) -> bool {
516        !self.has_partition_filter()
517            && self
518                .filters()
519                .iter()
520                .all(|filter| filter.semantic_type() == SemanticType::Tag)
521    }
522
523    /// Returns the format helper.
524    pub(crate) fn read_format(&self) -> &FlatReadFormat {
525        &self.base.read_format
526    }
527
528    /// Returns the reader builder.
529    pub(crate) fn reader_builder(&self) -> &RowGroupReaderBuilder {
530        &self.reader_builder
531    }
532
533    /// Returns the helper to compat batches.
534    pub(crate) fn compat_batch(&self) -> Option<&FlatCompatBatch> {
535        self.base.compat_batch.as_ref()
536    }
537
538    /// Returns the helper to project batches.
539    pub(crate) fn compaction_projection_mapper(&self) -> Option<&CompactionProjectionMapper> {
540        self.base.compaction_projection_mapper.as_ref()
541    }
542
543    /// Sets the compat helper to the context.
544    pub(crate) fn set_compat_batch(&mut self, compat: Option<FlatCompatBatch>) {
545        self.base.compat_batch = compat;
546    }
547
548    /// Filters the input RecordBatch by the pushed down predicate and returns RecordBatch.
549    /// If a partition expr filter is configured, it is also applied.
550    /// Physical filter exprs are not evaluated here; they are only applied during prefiltering.
551    pub(crate) fn precise_filter_flat(
552        &self,
553        input: RecordBatch,
554        skip_fields: bool,
555        skip_tags: bool,
556    ) -> Result<Option<RecordBatch>> {
557        self.base.precise_filter_flat(input, skip_fields, skip_tags)
558    }
559
560    pub(crate) fn pre_filter_mode(&self) -> PreFilterMode {
561        self.base.pre_filter_mode
562    }
563
564    //// Decodes parquet metadata and finds if row group contains delete op.
565    pub(crate) fn contains_delete(&self, row_group_index: usize) -> Result<bool> {
566        let metadata = self.reader_builder.parquet_metadata();
567        row_group_contains_delete(metadata, row_group_index, self.reader_builder.file_path())
568    }
569
570    /// Creates a [RowGroupBuildContext] for building row group readers with prefiltering.
571    pub(crate) fn build_context<'a>(
572        &'a self,
573        row_group_idx: usize,
574        row_selection: Option<RowSelection>,
575        fetch_metrics: Option<&'a ParquetFetchMetrics>,
576    ) -> RowGroupBuildContext<'a> {
577        RowGroupBuildContext {
578            row_group_idx,
579            row_selection,
580            fetch_metrics,
581        }
582    }
583
584    /// Returns the estimated memory size of this context.
585    /// Mainly accounts for the parquet metadata size.
586    pub(crate) fn memory_size(&self) -> usize {
587        self.reader_builder.parquet_metadata_size()
588    }
589}
590
591/// Mode to pre-filter columns in a range.
592#[derive(Debug, Clone, Copy, PartialEq, Eq)]
593pub enum PreFilterMode {
594    /// Filters all columns.
595    All,
596    /// Always skip fields.
597    SkipFields,
598}
599
600impl PreFilterMode {
601    pub(crate) fn skip_fields(self) -> bool {
602        matches!(self, Self::SkipFields)
603    }
604}
605
606/// Context for partition expression filtering.
607pub(crate) struct PartitionFilterContext {
608    pub(crate) region_partition_physical_expr: Arc<dyn PhysicalExpr>,
609    /// Schema containing only columns referenced by the partition expression.
610    /// This is used to build a minimal RecordBatch for partition filter evaluation.
611    pub(crate) partition_schema: Arc<Schema>,
612}
613
614/// Common fields for a range to read and filter batches.
615pub(crate) struct RangeBase {
616    /// Filters pushed down.
617    pub(crate) filters: Vec<SimpleFilterContext>,
618    /// Dynamic filter physical exprs.
619    pub(crate) dyn_filters: Vec<Arc<DynamicFilterPhysicalExpr>>,
620    /// Helper to read the SST.
621    pub(crate) read_format: FlatReadFormat,
622    pub(crate) expected_metadata: Option<RegionMetadataRef>,
623    /// Schema used for pruning with dynamic filters.
624    pub(crate) prune_schema: Arc<Schema>,
625    /// Decoder for primary keys
626    pub(crate) codec: Arc<dyn PrimaryKeyCodec>,
627    /// Optional helper to compat batches.
628    pub(crate) compat_batch: Option<FlatCompatBatch>,
629    /// Optional helper to project batches.
630    pub(crate) compaction_projection_mapper: Option<CompactionProjectionMapper>,
631    /// Mode to pre-filter columns.
632    pub(crate) pre_filter_mode: PreFilterMode,
633    /// Partition filter.
634    pub(crate) partition_filter: Option<PartitionFilterContext>,
635}
636
637pub(crate) struct TagDecodeState {
638    decoded_pks: Option<DecodedPrimaryKeys>,
639    decoded_tag_cache: HashMap<ColumnId, ArrayRef>,
640}
641
642impl TagDecodeState {
643    pub(crate) fn new() -> Self {
644        Self {
645            decoded_pks: None,
646            decoded_tag_cache: HashMap::new(),
647        }
648    }
649}
650
651impl RangeBase {
652    /// Filters the input RecordBatch by the pushed down predicate and returns RecordBatch.
653    ///
654    /// It assumes all necessary tags are already decoded from the primary key.
655    ///
656    /// # Arguments
657    /// * `input` - The RecordBatch to filter
658    /// * `skip_fields` - Whether to skip field filters based on PreFilterMode
659    /// * `skip_tags` - Whether to skip tag filters that were applied in an earlier phase
660    pub(crate) fn precise_filter_flat(
661        &self,
662        input: RecordBatch,
663        skip_fields: bool,
664        skip_tags: bool,
665    ) -> Result<Option<RecordBatch>> {
666        let mut tag_decode_state = TagDecodeState::new();
667        let mask =
668            self.compute_filter_mask_flat(&input, skip_fields, skip_tags, &mut tag_decode_state)?;
669
670        // If mask is None, the entire batch is filtered out
671        let Some(mut mask) = mask else {
672            return Ok(None);
673        };
674
675        // Apply partition filter
676        if let Some(partition_filter) = &self.partition_filter {
677            let record_batch = self.project_record_batch_for_pruning_flat(
678                &input,
679                &partition_filter.partition_schema,
680                &mut tag_decode_state,
681            )?;
682            let partition_mask = self.evaluate_partition_filter(&record_batch, partition_filter)?;
683            mask = mask.bitand(&partition_mask);
684        }
685
686        let num_selected = mask.count_set_bits();
687        if num_selected == 0 {
688            return Ok(None);
689        }
690        if num_selected == input.num_rows() {
691            // Nothing was filtered out, e.g. all filters were skipped by
692            // `skip_fields`/`skip_tags`. Avoid copying the whole batch.
693            return Ok(Some(input));
694        }
695
696        let filtered_batch =
697            datatypes::arrow::compute::filter_record_batch(&input, &BooleanArray::from(mask))
698                .context(ComputeArrowSnafu)?;
699
700        if filtered_batch.num_rows() > 0 {
701            Ok(Some(filtered_batch))
702        } else {
703            Ok(None)
704        }
705    }
706
707    /// Computes the filter mask for the input RecordBatch based on pushed down predicates.
708    /// If a partition expr filter is configured, it is applied later in `precise_filter_flat` but **NOT** in this function.
709    /// Physical filter exprs are excluded here and only apply during prefiltering.
710    ///
711    /// Returns `None` if the entire batch is filtered out, otherwise returns the boolean mask.
712    ///
713    /// # Arguments
714    /// * `input` - The RecordBatch to compute mask for
715    /// * `skip_fields` - Whether to skip field filters based on PreFilterMode
716    /// * `skip_tags` - Whether to skip tag filters that were applied in an earlier phase
717    pub(crate) fn compute_filter_mask_flat(
718        &self,
719        input: &RecordBatch,
720        skip_fields: bool,
721        skip_tags: bool,
722        tag_decode_state: &mut TagDecodeState,
723    ) -> Result<Option<BooleanBuffer>> {
724        let metadata = self.read_format.metadata();
725        // A pruned column rejects the batch without decoding any tag payloads.
726        if self
727            .filters
728            .iter()
729            .any(|ctx| matches!(ctx.filter(), MaybeFilter::Pruned))
730        {
731            return Ok(None);
732        }
733        let filter_tags = self
734            .filters
735            .iter()
736            .filter(|ctx| {
737                !skip_tags
738                    && ctx.semantic_type() == SemanticType::Tag
739                    && matches!(ctx.filter(), MaybeFilter::Filter(_))
740            })
741            .map(|ctx| ctx.column_id());
742        let partition_tags = self.partition_filter.iter().flat_map(|filter| {
743            filter
744                .partition_schema
745                .column_schemas()
746                .iter()
747                .filter_map(|column| {
748                    metadata
749                        .column_by_name(&column.name)
750                        .map(|column| column.column_id)
751                })
752        });
753        // Partition predicates are evaluated later, but extracting their tags
754        // together with the simple filters avoids a second scan of each key.
755        self.cache_sparse_tag_columns(input, filter_tags.chain(partition_tags), tag_decode_state)?;
756
757        let mut mask = BooleanBuffer::new_set(input.num_rows());
758
759        // Run filter one by one and combine them result
760        for filter_ctx in &self.filters {
761            let filter = match filter_ctx.filter() {
762                MaybeFilter::Filter(f) => f,
763                // Column matches.
764                MaybeFilter::Matched => continue,
765                // Column doesn't match, filter the entire batch.
766                MaybeFilter::Pruned => return Ok(None),
767            };
768
769            // Skip field filters if skip_fields is true
770            if skip_fields && filter_ctx.semantic_type() == SemanticType::Field {
771                continue;
772            }
773            if skip_tags && filter_ctx.semantic_type() == SemanticType::Tag {
774                continue;
775            }
776
777            // Get the column directly by its projected index.
778            // If the column is missing and it's not a tag/time column, this filter is skipped.
779            // Assumes the projection indices align with the input batch schema.
780            let column_idx = self
781                .read_format
782                .projected_index_by_id(filter_ctx.column_id());
783            if let Some(idx) = column_idx {
784                let column = &input.columns().get(idx).unwrap();
785                let result = filter.evaluate_array(column).context(RecordBatchSnafu)?;
786                mask = mask.bitand(&result);
787            } else if filter_ctx.semantic_type() == SemanticType::Tag {
788                // Column not found in projection, it may be a tag column.
789                let column_id = filter_ctx.column_id();
790
791                if let Some(tag_column) =
792                    self.maybe_decode_tag_column(metadata, column_id, input, tag_decode_state)?
793                {
794                    let result = filter
795                        .evaluate_array(&tag_column)
796                        .context(RecordBatchSnafu)?;
797                    mask = mask.bitand(&result);
798                }
799            } else if filter_ctx.semantic_type() == SemanticType::Timestamp {
800                let time_index_pos = time_index_column_index(input.num_columns());
801                let column = &input.columns()[time_index_pos];
802                let result = filter.evaluate_array(column).context(RecordBatchSnafu)?;
803                mask = mask.bitand(&result);
804            }
805            // Non-tag column not found in projection.
806        }
807
808        Ok(Some(mask))
809    }
810
811    /// Materializes missing sparse tags together, preserving the single-column fast path.
812    fn cache_sparse_tag_columns(
813        &self,
814        input: &RecordBatch,
815        column_ids: impl Iterator<Item = ColumnId>,
816        state: &mut TagDecodeState,
817    ) -> Result<()> {
818        if self.codec.encoding() != PrimaryKeyEncoding::Sparse {
819            return Ok(());
820        }
821        let metadata = self.read_format.metadata();
822        let mut seen = HashSet::new();
823        let columns: Vec<_> = column_ids
824            .filter(|id| {
825                metadata.primary_key_index(*id).is_some()
826                    && self.read_format.projected_index_by_id(*id).is_none()
827                    && !state.decoded_tag_cache.contains_key(id)
828                    && seen.insert(*id)
829            })
830            .filter_map(|id| {
831                metadata
832                    .column_by_id(id)
833                    .map(|column| (id, column.column_schema.data_type.clone()))
834            })
835            .collect();
836        if columns.is_empty() {
837            return Ok(());
838        }
839        let decoded = match &mut state.decoded_pks {
840            Some(decoded) => decoded,
841            None => state
842                .decoded_pks
843                .insert(decode_primary_keys(self.codec.as_ref(), input)?),
844        };
845        if let [(column_id, column_type)] = columns.as_slice() {
846            let array = decoded.get_tag_column(*column_id, None, column_type)?;
847            state.decoded_tag_cache.insert(*column_id, array);
848        } else {
849            let arrays = decoded.get_sparse_tag_columns(&columns)?;
850            state.decoded_tag_cache.extend(
851                columns
852                    .into_iter()
853                    .map(|(column_id, _)| column_id)
854                    .zip(arrays),
855            );
856        }
857        Ok(())
858    }
859
860    /// Returns the decoded tag column for `column_id`, or `None` if it's not a tag.
861    fn maybe_decode_tag_column(
862        &self,
863        metadata: &RegionMetadataRef,
864        column_id: ColumnId,
865        input: &RecordBatch,
866        tag_decode_state: &mut TagDecodeState,
867    ) -> Result<Option<ArrayRef>> {
868        let Some(pk_index) = metadata.primary_key_index(column_id) else {
869            return Ok(None);
870        };
871
872        if let Some(cached_column) = tag_decode_state.decoded_tag_cache.get(&column_id) {
873            return Ok(Some(cached_column.clone()));
874        }
875
876        if tag_decode_state.decoded_pks.is_none() {
877            tag_decode_state.decoded_pks = Some(decode_primary_keys(self.codec.as_ref(), input)?);
878        }
879
880        let pk_index = if self.codec.encoding() == PrimaryKeyEncoding::Sparse {
881            None
882        } else {
883            Some(pk_index)
884        };
885        let Some(column_index) = metadata.column_index_by_id(column_id) else {
886            return Ok(None);
887        };
888        let Some(decoded) = tag_decode_state.decoded_pks.as_mut() else {
889            return Ok(None);
890        };
891
892        let column_metadata = &metadata.column_metadatas[column_index];
893        let tag_column = decoded.get_tag_column(
894            column_id,
895            pk_index,
896            &column_metadata.column_schema.data_type,
897        )?;
898        tag_decode_state
899            .decoded_tag_cache
900            .insert(column_id, tag_column.clone());
901
902        Ok(Some(tag_column))
903    }
904
905    /// Evaluates the partition filter against the input `RecordBatch`.
906    fn evaluate_partition_filter(
907        &self,
908        record_batch: &RecordBatch,
909        partition_filter: &PartitionFilterContext,
910    ) -> Result<BooleanBuffer> {
911        let columnar_value = partition_filter
912            .region_partition_physical_expr
913            .evaluate(record_batch)
914            .context(EvalPartitionFilterSnafu)?;
915        let array = columnar_value
916            .into_array(record_batch.num_rows())
917            .context(EvalPartitionFilterSnafu)?;
918        let boolean_array =
919            array
920                .as_any()
921                .downcast_ref::<BooleanArray>()
922                .context(UnexpectedSnafu {
923                    reason: "Failed to downcast to BooleanArray".to_string(),
924                })?;
925
926        // also need to consider nulls in the partition filter result. If a value is null, it should be treated as false (filtered out).
927        let mut mask = boolean_array.values().clone();
928        if let Some(nulls) = boolean_array.nulls() {
929            mask = mask.bitand(nulls.inner());
930        }
931
932        Ok(mask)
933    }
934
935    /// Projects the input `RecordBatch` to match the given schema.
936    ///
937    /// This is used for partition expression evaluation. The schema should only contain
938    /// the columns referenced by the partition expression to minimize overhead.
939    fn project_record_batch_for_pruning_flat(
940        &self,
941        input: &RecordBatch,
942        schema: &Arc<Schema>,
943        tag_decode_state: &mut TagDecodeState,
944    ) -> Result<RecordBatch> {
945        let arrow_schema = schema.arrow_schema();
946        let mut columns = Vec::with_capacity(arrow_schema.fields().len());
947
948        let metadata = self.read_format.metadata();
949
950        for field in arrow_schema.fields() {
951            let column_id = metadata.column_by_name(field.name()).map(|c| c.column_id);
952
953            let Some(column_id) = column_id else {
954                return UnexpectedSnafu {
955                    reason: format!(
956                        "Partition pruning schema expects column '{}' but it is missing in \
957                         region metadata",
958                        field.name()
959                    ),
960                }
961                .fail();
962            };
963
964            if let Some(idx) = self.read_format.projected_index_by_id(column_id) {
965                columns.push(input.column(idx).clone());
966                continue;
967            }
968
969            if metadata.time_index_column().column_id == column_id {
970                let time_index_pos = time_index_column_index(input.num_columns());
971                columns.push(input.column(time_index_pos).clone());
972                continue;
973            }
974
975            if let Some(tag_column) =
976                self.maybe_decode_tag_column(metadata, column_id, input, tag_decode_state)?
977            {
978                columns.push(tag_column);
979                continue;
980            }
981
982            return UnexpectedSnafu {
983                reason: format!(
984                    "Partition pruning schema expects column '{}' (id {}) but it is not \
985                     present in projected record batch",
986                    field.name(),
987                    column_id
988                ),
989            }
990            .fail();
991        }
992
993        RecordBatch::try_new(arrow_schema.clone(), columns).context(NewRecordBatchSnafu)
994    }
995}
996
997#[cfg(test)]
998mod tests {
999    use std::sync::Arc;
1000
1001    use datafusion_expr::{col, lit};
1002    use datatypes::arrow::array::{
1003        BinaryDictionaryBuilder, TimestampMillisecondArray, UInt8Array, UInt64Array,
1004    };
1005    use datatypes::arrow::datatypes::UInt32Type;
1006    use datatypes::prelude::ConcreteDataType;
1007    use datatypes::schema::ColumnSchema;
1008    use datatypes::value::Value;
1009    use mito_codec::row_converter::SparsePrimaryKeyCodec;
1010    use parquet::arrow::arrow_reader::RowSelector;
1011    use partition::expr::col as partition_col;
1012
1013    use super::*;
1014    use crate::read::read_columns::ReadColumns;
1015    use crate::sst::parquet::flat_format::FlatReadFormat;
1016    use crate::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema};
1017    use crate::test_util::sst_util::{
1018        new_record_batch_with_custom_sequence, sst_region_metadata,
1019        sst_region_metadata_with_encoding,
1020    };
1021
1022    fn new_test_range_base(filters: Vec<SimpleFilterContext>) -> RangeBase {
1023        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1024
1025        let read_format = FlatReadFormat::new(
1026            metadata.clone(),
1027            ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
1028            None,
1029            "test",
1030            true,
1031        )
1032        .unwrap();
1033
1034        RangeBase {
1035            filters,
1036            dyn_filters: vec![],
1037            read_format,
1038            expected_metadata: None,
1039            prune_schema: metadata.schema.clone(),
1040            codec: mito_codec::row_converter::build_primary_key_codec(metadata.as_ref()),
1041            compat_batch: None,
1042            compaction_projection_mapper: None,
1043            pre_filter_mode: PreFilterMode::All,
1044            partition_filter: None,
1045        }
1046    }
1047
1048    #[test]
1049    fn test_compute_filter_mask_flat_applies_remaining_simple_filters() {
1050        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1051        let filters = vec![
1052            SimpleFilterContext::new_opt(&metadata, None, &col("tag_0").eq(lit("a"))).unwrap(),
1053            SimpleFilterContext::new_opt(&metadata, None, &col("field_0").gt(lit(1_u64))).unwrap(),
1054        ];
1055        let base = new_test_range_base(filters);
1056        let batch = new_record_batch_with_custom_sequence(&["b", "x"], 0, 4, 1);
1057
1058        let mask = base
1059            .compute_filter_mask_flat(&batch, false, false, &mut TagDecodeState::new())
1060            .unwrap()
1061            .unwrap();
1062        assert_eq!(mask.count_set_bits(), 0);
1063
1064        let mask = base
1065            .compute_filter_mask_flat(&batch, false, true, &mut TagDecodeState::new())
1066            .unwrap()
1067            .unwrap();
1068        assert_eq!(mask.count_set_bits(), 2);
1069    }
1070
1071    #[test]
1072    fn test_precise_filter_flat_returns_input_when_nothing_is_filtered() {
1073        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1074        let tag_filter =
1075            SimpleFilterContext::new_opt(&metadata, None, &col("tag_0").eq(lit("z"))).unwrap();
1076        let base = new_test_range_base(vec![tag_filter]);
1077        let batch = new_record_batch_with_custom_sequence(&["b", "x"], 0, 4, 1);
1078
1079        // The only filter is a tag filter, and it is skipped, so the batch must
1080        // come back untouched rather than being copied through `filter_record_batch`.
1081        let filtered = base
1082            .precise_filter_flat(batch.clone(), false, true)
1083            .unwrap()
1084            .unwrap();
1085        assert_eq!(batch, filtered);
1086    }
1087
1088    #[test]
1089    fn test_sparse_filter_and_partition_tags_preserve_skip_modes() {
1090        let metadata = Arc::new(sst_region_metadata_with_encoding(
1091            PrimaryKeyEncoding::Sparse,
1092        ));
1093        let codec = SparsePrimaryKeyCodec::schemaless();
1094        let long_tag = "x".repeat(80);
1095        let mut pk_builder = BinaryDictionaryBuilder::<UInt32Type>::new();
1096        for (series, tag0, tag1) in [
1097            (0, Some(long_tag.as_str()), "east"),
1098            (1, Some(""), "east"),
1099            (2, Some("b"), "west"),
1100            (3, None, "east"),
1101            (4, Some("c"), "east"),
1102            (0, Some(long_tag.as_str()), "east"),
1103        ] {
1104            let mut pk = Vec::new();
1105            codec.encode_internal(1, series, &mut pk).unwrap();
1106            let tags = tag0
1107                .map(|tag| (0, tag.as_bytes()))
1108                .into_iter()
1109                .chain([(1, tag1.as_bytes())]);
1110            codec.encode_raw_tag_value(tags, &mut pk).unwrap();
1111            pk_builder.append(pk).unwrap();
1112        }
1113        let raw = RecordBatch::try_new(
1114            to_flat_sst_arrow_schema(
1115                &metadata,
1116                &FlatSchemaOptions::from_encoding(PrimaryKeyEncoding::Sparse),
1117            ),
1118            vec![
1119                Arc::new(UInt64Array::from_iter_values(0..6)),
1120                Arc::new(TimestampMillisecondArray::from_iter_values(0..6)),
1121                Arc::new(pk_builder.finish()),
1122                Arc::new(UInt64Array::from(vec![1; 6])),
1123                Arc::new(UInt8Array::from(vec![0; 6])),
1124            ],
1125        )
1126        .unwrap();
1127        let filters: Vec<_> = [
1128            col("tag_0").gt(lit("")),
1129            col("tag_0").lt(lit("z")),
1130            col("field_0").gt(lit(1_u64)),
1131        ]
1132        .iter()
1133        .map(|expr| SimpleFilterContext::new_opt(&metadata, None, expr).unwrap())
1134        .collect();
1135        // Match build_partition_filter's dictionary schema for string tags.
1136        let partition_schema = Arc::new(Schema::new(vec![ColumnSchema::new(
1137            "tag_1",
1138            ConcreteDataType::dictionary_datatype(
1139                ConcreteDataType::uint32_datatype(),
1140                ConcreteDataType::string_datatype(),
1141            ),
1142            true,
1143        )]));
1144        let partition_expr = partition_col("tag_1")
1145            .eq(Value::String("east".into()))
1146            .try_as_physical_expr(partition_schema.arrow_schema())
1147            .unwrap();
1148
1149        // The partition tag is always encoded; tag_0 may already be materialized.
1150        for materialized in [false, true] {
1151            let read_format = FlatReadFormat::new(
1152                metadata.clone(),
1153                ReadColumns::new([0, 2, 3]),
1154                None,
1155                "test",
1156                !materialized,
1157            )
1158            .unwrap();
1159            let batch = read_format.convert_batch(raw.clone(), None).unwrap();
1160            let base = RangeBase {
1161                filters: filters.clone(),
1162                dyn_filters: vec![],
1163                read_format,
1164                expected_metadata: None,
1165                prune_schema: metadata.schema.clone(),
1166                codec: Arc::new(codec.clone()),
1167                compat_batch: None,
1168                compaction_projection_mapper: None,
1169                pre_filter_mode: PreFilterMode::All,
1170                partition_filter: Some(PartitionFilterContext {
1171                    partition_schema: partition_schema.clone(),
1172                    region_partition_physical_expr: partition_expr.clone(),
1173                }),
1174            };
1175            for (skip_fields, skip_tags, expected) in [
1176                (false, false, vec![4, 5]),
1177                (true, false, vec![0, 4, 5]),
1178                (false, true, vec![3, 4, 5]),
1179                (true, true, vec![0, 1, 3, 4, 5]),
1180            ] {
1181                let output = base
1182                    .precise_filter_flat(batch.clone(), skip_fields, skip_tags)
1183                    .unwrap()
1184                    .unwrap();
1185                let timestamps = output
1186                    .column(time_index_column_index(output.num_columns()))
1187                    .as_any()
1188                    .downcast_ref::<TimestampMillisecondArray>()
1189                    .unwrap();
1190                assert_eq!(
1191                    timestamps.values().as_ref(),
1192                    expected.as_slice(),
1193                    "materialized={materialized}, skip_fields={skip_fields}, skip_tags={skip_tags}"
1194                );
1195            }
1196        }
1197    }
1198
1199    #[test]
1200    fn test_compute_filter_mask_flat_does_not_postfilter_physical_filters() {
1201        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1202        let read_format = FlatReadFormat::new(
1203            metadata.clone(),
1204            ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
1205            None,
1206            "test",
1207            true,
1208        )
1209        .unwrap();
1210        let physical_filter = crate::sst::parquet::reader::PhysicalFilterContext::new_opt(
1211            &metadata,
1212            None,
1213            &read_format,
1214            &col("field_0").in_list(vec![lit(1_u64), lit(2_u64)], false),
1215        );
1216        assert!(physical_filter.is_some());
1217        let base = new_test_range_base(vec![]);
1218        let batch = new_record_batch_with_custom_sequence(&["b", "x"], 0, 4, 1);
1219
1220        let mask = base
1221            .compute_filter_mask_flat(&batch, false, false, &mut TagDecodeState::new())
1222            .unwrap()
1223            .unwrap();
1224        assert_eq!(mask.count_set_bits(), 4);
1225    }
1226
1227    #[test]
1228    fn test_precise_filter_flat_applies_partition_filter_when_skipping_tags() {
1229        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1230        let tag_filter =
1231            SimpleFilterContext::new_opt(&metadata, None, &col("tag_0").eq(lit("z"))).unwrap();
1232        let mut base = new_test_range_base(vec![tag_filter]);
1233        let batch = new_record_batch_with_custom_sequence(&["b", "x"], 0, 4, 1);
1234
1235        let batch_schema = batch.schema();
1236        let tag_field = batch_schema.field(0);
1237        let partition_schema = Arc::new(Schema::new(vec![ColumnSchema::new(
1238            "tag_0".to_string(),
1239            ConcreteDataType::from_arrow_type(tag_field.data_type()),
1240            tag_field.is_nullable(),
1241        )]));
1242        let partition_expr = partition_col("tag_0")
1243            .gt_eq(Value::String("a".into()))
1244            .and(partition_col("tag_0").lt(Value::String("c".into())));
1245        base.partition_filter = Some(PartitionFilterContext {
1246            region_partition_physical_expr: partition_expr
1247                .try_as_physical_expr(partition_schema.arrow_schema())
1248                .unwrap(),
1249            partition_schema,
1250        });
1251
1252        let filtered = base
1253            .precise_filter_flat(batch, false, true)
1254            .unwrap()
1255            .unwrap();
1256        assert_eq!(filtered.num_rows(), 4);
1257
1258        let out_of_partition = new_record_batch_with_custom_sequence(&["z", "x"], 0, 4, 1);
1259        assert!(
1260            base.precise_filter_flat(out_of_partition, false, true)
1261                .unwrap()
1262                .is_none()
1263        );
1264    }
1265
1266    #[test]
1267    fn test_refine_row_range_selection() {
1268        let original = Some(RowSelection::from(vec![
1269            RowSelector::skip(2),
1270            RowSelector::select(3),
1271            RowSelector::skip(1),
1272            RowSelector::select(2),
1273        ]));
1274        let selected = refine_row_range_selection(vec![0..3, 4..7, 9..10], 10, &original)
1275            .unwrap()
1276            .unwrap();
1277        assert_eq!(
1278            selected,
1279            RowSelection::from(vec![
1280                RowSelector::skip(2),
1281                RowSelector::select(1),
1282                RowSelector::skip(1),
1283                RowSelector::select(1),
1284                RowSelector::skip(1),
1285                RowSelector::select(1),
1286                RowSelector::skip(1),
1287            ])
1288        );
1289        assert!(
1290            refine_row_range_selection(vec![0..2, 8..10], 10, &original)
1291                .unwrap()
1292                .is_none()
1293        );
1294        assert!(
1295            refine_row_range_selection(vec![], 10, &None)
1296                .unwrap()
1297                .is_none()
1298        );
1299        assert_eq!(
1300            refine_row_range_selection(vec![2..4, 4..6], 10, &None)
1301                .unwrap()
1302                .unwrap(),
1303            RowSelection::from(vec![RowSelector::skip(2), RowSelector::select(4)])
1304        );
1305        for ranges in [
1306            std::iter::once(0..11).collect(),
1307            std::iter::once(11..12).collect(),
1308            std::iter::once(3..3).collect(),
1309            vec![4..6, 5..7],
1310            vec![4..6, 0..2],
1311        ] {
1312            assert!(refine_row_range_selection(ranges, 10, &None).is_err());
1313        }
1314    }
1315
1316    #[test]
1317    fn test_refine_primary_key_selection_intersects_original_selection() {
1318        let original = Some(RowSelection::from(vec![
1319            RowSelector::skip(2),
1320            RowSelector::select(3),
1321            RowSelector::skip(1),
1322            RowSelector::select(2),
1323        ]));
1324        let mask = BooleanArray::from(vec![true, false, true, false, true]);
1325        let actual = refine_primary_key_selection(&[mask], &original).unwrap();
1326        let expected = RowSelection::from(vec![
1327            RowSelector::skip(2),
1328            RowSelector::select(1),
1329            RowSelector::skip(1),
1330            RowSelector::select(1),
1331            RowSelector::skip(2),
1332            RowSelector::select(1),
1333        ]);
1334        assert_eq!(actual, expected);
1335    }
1336}