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 { 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>> { + let file_size = data.len() as u64; + let calls = Arc::new(Mutex::new(Vec::new())); + let _: Vec = 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 { 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| { + 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 { + 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::>(); + let values = ids.iter().map(|value| value * 10).collect::>(); + 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::( + 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();