From b09c34ccab1884d5a0d4f47e0d0440d25a43de16 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Thu, 23 Jul 2026 15:13:45 +0800 Subject: [PATCH] perf(parquet): prune row selection with OffsetIndex --- crates/paimon/src/arrow/format/parquet.rs | 139 ++++++++++++++++++++-- 1 file changed, 128 insertions(+), 11 deletions(-) diff --git a/crates/paimon/src/arrow/format/parquet.rs b/crates/paimon/src/arrow/format/parquet.rs index 4d4d566cf..c851299f3 100644 --- a/crates/paimon/src/arrow/format/parquet.rs +++ b/crates/paimon/src/arrow/format/parquet.rs @@ -310,17 +310,14 @@ impl FormatFileReader for ParquetFormatReader { }; let row_filter_factory = predicates.and_then(|fp| fp.row_filter_factory.as_deref()); - // Only load the Parquet page index (ColumnIndex + OffsetIndex) when a - // predicate can use it for page-level pruning — matching Java Paimon, - // which gets page-level skipping for free via parquet-mr's - // `readNextFilteredRowGroup`. arrow-rs does not do this automatically, so - // we build the RowSelection ourselves below. Without a predicate the - // index is pure overhead (an extra metadata read with no benefit), so we - // skip it. `Optional` lets files without a page index fall through to - // row-group-level pruning instead of erroring. + // Predicates need both indexes for page-stat pruning. Row selection only + // needs OffsetIndex so arrow-rs can avoid fetching unselected pages. let mut arrow_options = ArrowReaderOptions::new(); if !preds.is_empty() { - arrow_options = arrow_options.with_page_index_policy(PageIndexPolicy::Optional); + arrow_options = arrow_options.with_column_index_policy(PageIndexPolicy::Optional); + } + if !preds.is_empty() || row_selection.is_some() { + arrow_options = arrow_options.with_offset_index_policy(PageIndexPolicy::Optional); } let mut batch_stream_builder = ParquetRecordBatchStreamBuilder::new_with_options(arrow_file_reader, arrow_options) @@ -1988,9 +1985,10 @@ mod tests { use crate::arrow::{build_target_arrow_schema, variant_arrow_type}; use crate::io::FileIOBuilder; use crate::spec::{ - BigIntType, DataField, DataType, Datum, IntType, MapType, PredicateBuilder, VarCharType, - VariantType, + ArrayType, BigIntType, DataField, DataType, Datum, IntType, MapType, PredicateBuilder, + VarCharType, VariantType, }; + use crate::table::RowRange; use crate::variant::GenericVariant; use arrow_array::{ Array, BinaryArray, Int32Array, Int64Array, MapArray, RecordBatch, StringArray, StructArray, @@ -2835,6 +2833,125 @@ mod tests { buf } + #[derive(Clone)] + struct TrackingFileRead { + data: Bytes, + ranges: Arc>>>, + } + + impl TrackingFileRead { + fn new(data: Bytes) -> Self { + Self { + data, + ranges: Arc::new(std::sync::Mutex::new(Vec::new())), + } + } + + fn bytes_read(&self) -> u64 { + self.ranges + .lock() + .unwrap() + .iter() + .map(|range| range.end - range.start) + .sum() + } + + fn reset(&self) { + self.ranges.lock().unwrap().clear(); + } + } + + #[async_trait::async_trait] + impl crate::io::FileRead for TrackingFileRead { + async fn read(&self, range: std::ops::Range) -> crate::Result { + self.ranges.lock().unwrap().push(range.clone()); + Ok(self.data.slice(range.start as usize..range.end as usize)) + } + } + + async fn write_nested_multi_page_parquet() -> Vec { + use arrow_array::builder::{Int32Builder, ListBuilder}; + + const ROWS: usize = 1024; + const VALUES_PER_ROW: usize = 1024; + + let element = Arc::new(ArrowField::new("element", ArrowDataType::Int32, false)); + let mut values = ListBuilder::new(Int32Builder::new()).with_field(element.clone()); + for row in 0..ROWS { + for value in 0..VALUES_PER_ROW { + values + .values() + .append_value((row * VALUES_PER_ROW + value) as i32); + } + values.append(true); + } + + let schema = Arc::new(ArrowSchema::new(vec![ArrowField::new( + "items", + ArrowDataType::List(element), + false, + )])); + let batch = RecordBatch::try_new(schema.clone(), vec![Arc::new(values.finish())]).unwrap(); + let props = parquet::file::properties::WriterProperties::builder() + .set_data_page_row_count_limit(128) + .set_write_batch_size(128) + .set_max_row_group_row_count(Some(ROWS)) + .set_dictionary_enabled(false) + .build(); + + let mut buf = Vec::new(); + let mut writer = AsyncArrowWriter::try_new(&mut buf, schema, Some(props)).unwrap(); + writer.write(&batch).await.unwrap(); + writer.close().await.unwrap(); + buf + } + + async fn read_nested_rows(data: Bytes, row_selection: Option>) -> (usize, u64) { + let file_size = data.len() as u64; + let file_read = TrackingFileRead::new(data); + let tracker = file_read.clone(); + let fields = vec![DataField::new( + 0, + "items".to_string(), + DataType::Array(ArrayType::new(DataType::Int(IntType::new()))), + )]; + let stream = ParquetFormatReader + .read_batch_stream( + Box::new(file_read), + file_size, + &fields, + None, + Some(128), + row_selection, + ) + .await + .unwrap(); + tracker.reset(); + let rows = stream + .try_fold( + 0usize, + |rows, batch| async move { Ok(rows + batch.num_rows()) }, + ) + .await + .unwrap(); + (rows, tracker.bytes_read()) + } + + #[tokio::test] + async fn test_row_selection_prunes_nested_page_io() { + let data = Bytes::from(write_nested_multi_page_parquet().await); + let (all_rows, all_bytes) = read_nested_rows(data.clone(), None).await; + let (selected_rows, selected_bytes) = + read_nested_rows(data, Some(vec![RowRange::new(0, 0)])).await; + + assert_eq!(all_rows, 1024); + assert_eq!(selected_rows, 1); + assert!( + selected_bytes * 2 < all_bytes, + "selected read used {selected_bytes} bytes; full read used {all_bytes} bytes" + ); + } + /// Parse metadata from in-memory parquet bytes, optionally loading the page /// index — mirrors what the reader does via `with_page_index_policy`. fn load_metadata_with_page_index(