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 5c17d26c perf(parquet): prune row selection with OffsetIndex (#597)
5c17d26c is described below
commit 5c17d26c8808b6e7c9b9a44bf10e09c3e32ef40a
Author: XiaoHongbo <[email protected]>
AuthorDate: Thu Jul 23 17:12:45 2026 +0800
perf(parquet): prune row selection with OffsetIndex (#597)
---
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 4d4d566c..c851299f 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<std::sync::Mutex<Vec<std::ops::Range<u64>>>>,
+ }
+
+ 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<u64>) ->
crate::Result<Bytes> {
+ 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<u8> {
+ 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<Vec<RowRange>>) -> (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(