This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new 2cdfec8d perf(arrow): prune IN predicates with row-group stats (#705)
2cdfec8d is described below

commit 2cdfec8d3012dcab37f0fb41d1c49e2ccbab9a3e
Author: QuakeWang <[email protected]>
AuthorDate: Sat Aug 15 16:16:07 2026 +0800

    perf(arrow): prune IN predicates with row-group stats (#705)
---
 crates/paimon/src/arrow/format/mosaic.rs  | 107 ++++++++++++++++++++++++++++
 crates/paimon/src/arrow/format/parquet.rs | 111 +++++++++++++++++++++++++++++-
 2 files changed, 217 insertions(+), 1 deletion(-)

diff --git a/crates/paimon/src/arrow/format/mosaic.rs 
b/crates/paimon/src/arrow/format/mosaic.rs
index 304eadc9..1fdbbace 100644
--- a/crates/paimon/src/arrow/format/mosaic.rs
+++ b/crates/paimon/src/arrow/format/mosaic.rs
@@ -264,6 +264,10 @@ impl StatsAccessor for MosaicRowGroupStats<'_> {
     fn max_value(&self, index: usize, data_type: &PaimonDataType) -> 
Option<Datum> {
         mosaic_value_to_datum(self.column_stats(index)?.max.as_ref()?, 
data_type)
     }
+
+    fn supports_in_min_max_pruning(&self) -> bool {
+        true
+    }
 }
 
 impl MosaicRowGroupStats<'_> {
@@ -873,6 +877,32 @@ mod tests {
             .await
     }
 
+    async fn read_ranges_with_predicates(
+        data: Bytes,
+        read_fields: &[DataField],
+        predicates: &FilePredicates,
+    ) -> crate::Result<Vec<Range<u64>>> {
+        let file_size = data.len() as u64;
+        let calls = Arc::new(Mutex::new(Vec::new()));
+        let _: Vec<RecordBatch> = MosaicFormatReader
+            .read_batch_stream(
+                Box::new(TrackingFileRead {
+                    data,
+                    calls: Arc::clone(&calls),
+                }),
+                file_size,
+                read_fields,
+                Some(predicates),
+                None,
+                None,
+            )
+            .await?
+            .try_collect()
+            .await?;
+        let ranges = calls.lock().unwrap().clone();
+        Ok(ranges)
+    }
+
     fn collect_i32_column(batches: &[RecordBatch], column_index: usize) -> 
Vec<i32> {
         batches
             .iter()
@@ -1191,6 +1221,83 @@ mod tests {
         assert!(batches.is_empty());
     }
 
+    #[test]
+    fn test_row_group_in_min_max_pruning() {
+        let fields = vec![data_fields()[0].clone()];
+        let builder = PredicateBuilder::new(&fields);
+        let mapping = [Some(0)];
+        let stats = |min, max| {
+            [ColumnStats {
+                column_index: 0,
+                null_count: 0,
+                min,
+                max,
+            }]
+        };
+        let may_match = |stats: &[ColumnStats], literals: Vec<Datum>| {
+            let predicate = builder.is_in("id", literals).unwrap();
+            row_group_may_match(10, stats, &mapping, &[predicate], 
&fields).unwrap()
+        };
+
+        let valid_stats = stats(
+            Some(MosaicValue::Integer(10)),
+            Some(MosaicValue::Integer(20)),
+        );
+        assert!(!may_match(
+            &valid_stats,
+            vec![Datum::Int(1), Datum::Int(30)]
+        ));
+        assert!(may_match(&valid_stats, vec![Datum::Int(1), Datum::Int(15)]));
+        assert!(may_match(&[], vec![Datum::Int(1), Datum::Int(30)]));
+
+        let damaged_stats = stats(
+            Some(MosaicValue::Integer(20)),
+            Some(MosaicValue::Integer(10)),
+        );
+        assert!(may_match(
+            &damaged_stats,
+            vec![Datum::Int(1), Datum::Int(30)]
+        ));
+
+        let incomparable_stats = stats(
+            Some(MosaicValue::String(b"a".to_vec())),
+            Some(MosaicValue::String(b"z".to_vec())),
+        );
+        assert!(may_match(
+            &incomparable_stats,
+            vec![Datum::Int(1), Datum::Int(30)]
+        ));
+    }
+
+    #[tokio::test]
+    async fn test_in_pruning_skips_mosaic_column_reads() {
+        let fields = data_fields();
+        let projected = vec![fields[0].clone()];
+        let builder = PredicateBuilder::new(&fields);
+        let eq_predicates = predicate_file_predicates(
+            fields.clone(),
+            vec![builder.equal("id", Datum::Int(99)).unwrap()],
+        );
+        let in_predicates = predicate_file_predicates(
+            fields.clone(),
+            vec![builder
+                .is_in("id", vec![Datum::Int(99), Datum::Int(100)])
+                .unwrap()],
+        );
+        let data = multi_row_group_mosaic(vec!["id".to_string()]);
+
+        let eq_reads = read_ranges_with_predicates(data.clone(), &projected, 
&eq_predicates)
+            .await
+            .unwrap();
+        let in_reads = read_ranges_with_predicates(data, &projected, 
&in_predicates)
+            .await
+            .unwrap();
+        assert_eq!(
+            in_reads, eq_reads,
+            "an all-outside IN should not read row-group column data"
+        );
+    }
+
     #[tokio::test]
     async fn test_read_predicate_missing_stats_still_filters_rows() {
         let fields = data_fields();
diff --git a/crates/paimon/src/arrow/format/parquet.rs 
b/crates/paimon/src/arrow/format/parquet.rs
index 960828ec..76c76892 100644
--- a/crates/paimon/src/arrow/format/parquet.rs
+++ b/crates/paimon/src/arrow/format/parquet.rs
@@ -1018,6 +1018,10 @@ impl StatsAccessor for ParquetRowGroupStats<'_> {
             false,
         )
     }
+
+    fn supports_in_min_max_pruning(&self) -> bool {
+        true
+    }
 }
 
 fn build_predicate_row_selection(
@@ -2211,6 +2215,8 @@ mod tests {
     use arrow_buffer::{NullBuffer, OffsetBuffer, ScalarBuffer};
     use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema 
as ArrowSchema};
     use futures::{StreamExt, TryStreamExt};
+    use parquet::file::properties::EnabledStatistics;
+    use parquet::file::statistics::Statistics as ParquetStatistics;
     use parquet::schema::{parser::parse_message_type, types::SchemaDescriptor};
     use std::collections::HashMap;
     use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering};
@@ -3283,6 +3289,110 @@ mod tests {
         }
     }
 
+    // -----------------------------------------------------------------------
+    // Row-group statistics pruning
+    // -----------------------------------------------------------------------
+
+    async fn write_multi_row_group_parquet(
+        row_group_rows: usize,
+        total_rows: i32,
+        statistics: EnabledStatistics,
+    ) -> Vec<u8> {
+        let schema = writer_arrow_schema();
+        let props = parquet::file::properties::WriterProperties::builder()
+            .set_max_row_group_row_count(Some(row_group_rows))
+            .set_statistics_enabled(statistics)
+            .build();
+        let mut buf = Vec::new();
+        let mut writer = AsyncArrowWriter::try_new(&mut buf, schema.clone(), 
Some(props)).unwrap();
+        let ids = (0..total_rows).collect::<Vec<_>>();
+        let values = ids.iter().map(|value| value * 10).collect::<Vec<_>>();
+        writer
+            .write(&writer_test_batch(&schema, ids, values))
+            .await
+            .unwrap();
+        writer.close().await.unwrap();
+        buf
+    }
+
+    #[tokio::test]
+    async fn test_row_group_selection_in_uses_min_max_without_page_index() {
+        let bytes = write_multi_row_group_parquet(10, 20, 
EnabledStatistics::Chunk).await;
+        let metadata = load_metadata_with_page_index(&bytes, false);
+        assert_eq!(metadata.row_groups().len(), 2);
+        assert!(metadata.column_index().is_none());
+        assert!(metadata.offset_index().is_none());
+        let fields = vec![int_field("id"), int_field("value")];
+
+        let outside = vec![id_leaf(
+            PredicateOperator::In,
+            vec![Datum::Int(-1), Datum::Int(30)],
+        )];
+        let selection =
+            super::build_predicate_row_selection(metadata.row_groups(), 
&outside, &fields)
+                .unwrap()
+                .expect("all row groups should be skipped");
+        assert_eq!(selection.row_count(), 0);
+
+        let overlapping = vec![id_leaf(
+            PredicateOperator::In,
+            vec![Datum::Int(5), Datum::Int(30)],
+        )];
+        let selection =
+            super::build_predicate_row_selection(metadata.row_groups(), 
&overlapping, &fields)
+                .unwrap()
+                .expect("only the first row group should be kept");
+        assert_eq!(selection.row_count(), 10);
+    }
+
+    #[tokio::test]
+    async fn test_row_group_selection_in_fails_open_on_unusable_stats() {
+        let fields = vec![int_field("id"), int_field("value")];
+        let predicates = vec![id_leaf(PredicateOperator::In, 
vec![Datum::Int(100)])];
+
+        let bytes = write_multi_row_group_parquet(10, 10, 
EnabledStatistics::None).await;
+        let metadata = load_metadata_with_page_index(&bytes, false);
+        let selection =
+            super::build_predicate_row_selection(metadata.row_groups(), 
&predicates, &fields)
+                .unwrap();
+        assert!(selection.is_none(), "missing stats must fail open");
+
+        let bytes = write_multi_row_group_parquet(10, 10, 
EnabledStatistics::Chunk).await;
+        let metadata = load_metadata_with_page_index(&bytes, false);
+        let mut damaged_row_group = metadata.row_groups()[0].clone();
+        let damaged_id_column = damaged_row_group
+            .column(0)
+            .clone()
+            .into_builder()
+            .set_statistics(ParquetStatistics::new::<i32>(
+                Some(20),
+                Some(10),
+                None,
+                Some(0),
+                false,
+            ))
+            .build()
+            .unwrap();
+        damaged_row_group.columns_mut()[0] = damaged_id_column;
+        let selection =
+            super::build_predicate_row_selection(&[damaged_row_group], 
&predicates, &fields)
+                .unwrap();
+        assert!(selection.is_none(), "inverted stats must fail open");
+
+        let varchar_type = DataType::VarChar(VarCharType::new(20).unwrap());
+        let varchar_fields = vec![DataField::new(0, "id".to_string(), 
varchar_type)];
+        let varchar_predicates = vec![PredicateBuilder::new(&varchar_fields)
+            .is_in("id", vec![Datum::String("100".to_string())])
+            .unwrap()];
+        let selection = super::build_predicate_row_selection(
+            metadata.row_groups(),
+            &varchar_predicates,
+            &varchar_fields,
+        )
+        .unwrap();
+        assert!(selection.is_none(), "incomparable stats must fail open");
+    }
+
     // -----------------------------------------------------------------------
     // Page-index (ColumnIndex / OffsetIndex) pruning
     // -----------------------------------------------------------------------
@@ -3628,7 +3738,6 @@ mod tests {
         // (`ColumnIndexMetaData::NONE`) but still gets an OffsetIndex. Its
         // accessors panic rather than return None, so pruning must fail open
         // for that column instead of touching the index.
-        use parquet::file::properties::EnabledStatistics;
         use parquet::schema::types::ColumnPath;
 
         let schema = writer_arrow_schema();

Reply via email to