Skip to content

Commit 5c17d26

Browse files
perf(parquet): prune row selection with OffsetIndex (#597)
1 parent 9e7accf commit 5c17d26

1 file changed

Lines changed: 128 additions & 11 deletions

File tree

crates/paimon/src/arrow/format/parquet.rs

Lines changed: 128 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -310,17 +310,14 @@ impl FormatFileReader for ParquetFormatReader {
310310
};
311311
let row_filter_factory = predicates.and_then(|fp| fp.row_filter_factory.as_deref());
312312

313-
// Only load the Parquet page index (ColumnIndex + OffsetIndex) when a
314-
// predicate can use it for page-level pruning — matching Java Paimon,
315-
// which gets page-level skipping for free via parquet-mr's
316-
// `readNextFilteredRowGroup`. arrow-rs does not do this automatically, so
317-
// we build the RowSelection ourselves below. Without a predicate the
318-
// index is pure overhead (an extra metadata read with no benefit), so we
319-
// skip it. `Optional` lets files without a page index fall through to
320-
// row-group-level pruning instead of erroring.
313+
// Predicates need both indexes for page-stat pruning. Row selection only
314+
// needs OffsetIndex so arrow-rs can avoid fetching unselected pages.
321315
let mut arrow_options = ArrowReaderOptions::new();
322316
if !preds.is_empty() {
323-
arrow_options = arrow_options.with_page_index_policy(PageIndexPolicy::Optional);
317+
arrow_options = arrow_options.with_column_index_policy(PageIndexPolicy::Optional);
318+
}
319+
if !preds.is_empty() || row_selection.is_some() {
320+
arrow_options = arrow_options.with_offset_index_policy(PageIndexPolicy::Optional);
324321
}
325322
let mut batch_stream_builder =
326323
ParquetRecordBatchStreamBuilder::new_with_options(arrow_file_reader, arrow_options)
@@ -1988,9 +1985,10 @@ mod tests {
19881985
use crate::arrow::{build_target_arrow_schema, variant_arrow_type};
19891986
use crate::io::FileIOBuilder;
19901987
use crate::spec::{
1991-
BigIntType, DataField, DataType, Datum, IntType, MapType, PredicateBuilder, VarCharType,
1992-
VariantType,
1988+
ArrayType, BigIntType, DataField, DataType, Datum, IntType, MapType, PredicateBuilder,
1989+
VarCharType, VariantType,
19931990
};
1991+
use crate::table::RowRange;
19941992
use crate::variant::GenericVariant;
19951993
use arrow_array::{
19961994
Array, BinaryArray, Int32Array, Int64Array, MapArray, RecordBatch, StringArray, StructArray,
@@ -2835,6 +2833,125 @@ mod tests {
28352833
buf
28362834
}
28372835

2836+
#[derive(Clone)]
2837+
struct TrackingFileRead {
2838+
data: Bytes,
2839+
ranges: Arc<std::sync::Mutex<Vec<std::ops::Range<u64>>>>,
2840+
}
2841+
2842+
impl TrackingFileRead {
2843+
fn new(data: Bytes) -> Self {
2844+
Self {
2845+
data,
2846+
ranges: Arc::new(std::sync::Mutex::new(Vec::new())),
2847+
}
2848+
}
2849+
2850+
fn bytes_read(&self) -> u64 {
2851+
self.ranges
2852+
.lock()
2853+
.unwrap()
2854+
.iter()
2855+
.map(|range| range.end - range.start)
2856+
.sum()
2857+
}
2858+
2859+
fn reset(&self) {
2860+
self.ranges.lock().unwrap().clear();
2861+
}
2862+
}
2863+
2864+
#[async_trait::async_trait]
2865+
impl crate::io::FileRead for TrackingFileRead {
2866+
async fn read(&self, range: std::ops::Range<u64>) -> crate::Result<Bytes> {
2867+
self.ranges.lock().unwrap().push(range.clone());
2868+
Ok(self.data.slice(range.start as usize..range.end as usize))
2869+
}
2870+
}
2871+
2872+
async fn write_nested_multi_page_parquet() -> Vec<u8> {
2873+
use arrow_array::builder::{Int32Builder, ListBuilder};
2874+
2875+
const ROWS: usize = 1024;
2876+
const VALUES_PER_ROW: usize = 1024;
2877+
2878+
let element = Arc::new(ArrowField::new("element", ArrowDataType::Int32, false));
2879+
let mut values = ListBuilder::new(Int32Builder::new()).with_field(element.clone());
2880+
for row in 0..ROWS {
2881+
for value in 0..VALUES_PER_ROW {
2882+
values
2883+
.values()
2884+
.append_value((row * VALUES_PER_ROW + value) as i32);
2885+
}
2886+
values.append(true);
2887+
}
2888+
2889+
let schema = Arc::new(ArrowSchema::new(vec![ArrowField::new(
2890+
"items",
2891+
ArrowDataType::List(element),
2892+
false,
2893+
)]));
2894+
let batch = RecordBatch::try_new(schema.clone(), vec![Arc::new(values.finish())]).unwrap();
2895+
let props = parquet::file::properties::WriterProperties::builder()
2896+
.set_data_page_row_count_limit(128)
2897+
.set_write_batch_size(128)
2898+
.set_max_row_group_row_count(Some(ROWS))
2899+
.set_dictionary_enabled(false)
2900+
.build();
2901+
2902+
let mut buf = Vec::new();
2903+
let mut writer = AsyncArrowWriter::try_new(&mut buf, schema, Some(props)).unwrap();
2904+
writer.write(&batch).await.unwrap();
2905+
writer.close().await.unwrap();
2906+
buf
2907+
}
2908+
2909+
async fn read_nested_rows(data: Bytes, row_selection: Option<Vec<RowRange>>) -> (usize, u64) {
2910+
let file_size = data.len() as u64;
2911+
let file_read = TrackingFileRead::new(data);
2912+
let tracker = file_read.clone();
2913+
let fields = vec![DataField::new(
2914+
0,
2915+
"items".to_string(),
2916+
DataType::Array(ArrayType::new(DataType::Int(IntType::new()))),
2917+
)];
2918+
let stream = ParquetFormatReader
2919+
.read_batch_stream(
2920+
Box::new(file_read),
2921+
file_size,
2922+
&fields,
2923+
None,
2924+
Some(128),
2925+
row_selection,
2926+
)
2927+
.await
2928+
.unwrap();
2929+
tracker.reset();
2930+
let rows = stream
2931+
.try_fold(
2932+
0usize,
2933+
|rows, batch| async move { Ok(rows + batch.num_rows()) },
2934+
)
2935+
.await
2936+
.unwrap();
2937+
(rows, tracker.bytes_read())
2938+
}
2939+
2940+
#[tokio::test]
2941+
async fn test_row_selection_prunes_nested_page_io() {
2942+
let data = Bytes::from(write_nested_multi_page_parquet().await);
2943+
let (all_rows, all_bytes) = read_nested_rows(data.clone(), None).await;
2944+
let (selected_rows, selected_bytes) =
2945+
read_nested_rows(data, Some(vec![RowRange::new(0, 0)])).await;
2946+
2947+
assert_eq!(all_rows, 1024);
2948+
assert_eq!(selected_rows, 1);
2949+
assert!(
2950+
selected_bytes * 2 < all_bytes,
2951+
"selected read used {selected_bytes} bytes; full read used {all_bytes} bytes"
2952+
);
2953+
}
2954+
28382955
/// Parse metadata from in-memory parquet bytes, optionally loading the page
28392956
/// index — mirrors what the reader does via `with_page_index_policy`.
28402957
fn load_metadata_with_page_index(

0 commit comments

Comments
 (0)