This is an automated email from the ASF dual-hosted git repository.
alamb pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-rs.git
The following commit(s) were added to refs/heads/main by this push:
new 21f1156821 test(parquet): show push decoder needs a whole row group
before decoding a batch (#11218)
21f1156821 is described below
commit 21f115682123b9a0affe53471244fc7c23a88c39
Author: Andrew Lamb <[email protected]>
AuthorDate: Mon Sep 28 16:54:03 2026 -0400
test(parquet): show push decoder needs a whole row group before decoding a
batch (#11218)
# Which issue does this PR close?
- Related to https://github.com/apache/arrow-rs/issues/6946
# Rationale for this change
While discussing https://github.com/apache/arrow-rs/issues/6946 with
@adriangb, a question came up about whether `ParquetPushDecoder` would
decode a batch if only the pages needed for that batch had been pushed
or if it insists on having all pages for the row group first.
This PR adds a test that answers that question
For the current code, the ParquetPushDecoder requires that all ranges
for the row group are present
# What changes are included in this PR?
`test_decoder_first_page_only` test
# Are these changes tested?
This PR is only a test.
# Are there any user-facing changes?
No.
---------
Co-authored-by: Claude Fable 5.1 <[email protected]>
---
parquet/src/arrow/push_decoder/mod.rs | 85 ++++++++++++++++++++++++++++++++++-
1 file changed, 83 insertions(+), 2 deletions(-)
diff --git a/parquet/src/arrow/push_decoder/mod.rs
b/parquet/src/arrow/push_decoder/mod.rs
index 24561c33b4..96982b68ff 100644
--- a/parquet/src/arrow/push_decoder/mod.rs
+++ b/parquet/src/arrow/push_decoder/mod.rs
@@ -972,7 +972,7 @@ mod test {
};
use crate::arrow::{ArrowWriter, ProjectionMask};
use crate::errors::ParquetError;
- use crate::file::metadata::ParquetMetaDataPushDecoder;
+ use crate::file::metadata::{PageIndexPolicy, ParquetMetaDataPushDecoder};
use crate::file::properties::WriterProperties;
use arrow::compute::kernels::cmp::{gt, lt};
use arrow_array::cast::AsArray;
@@ -1144,6 +1144,74 @@ mod test {
expect_finished(decoder.try_decode());
}
+ /// Push only the pages needed for each batch of a row group rather than
+ /// the entire row group, and check whether the decoder can produce that
+ /// batch before the rest of the row group's pages have been pushed.
+ #[test]
+ fn test_decoder_first_pages_only() {
+ let metadata = test_file_parquet_metadata_with_offset_index();
+ let mut decoder =
ParquetPushDecoderBuilder::try_new_decoder(Arc::clone(&metadata))
+ .unwrap()
+ // Each data page has 100 rows, so a batch of 100 rows needs only
+ // the first data page of each column
+ .with_batch_size(100)
+ .build()
+ .unwrap();
+
+ // Row group 0: the decoder asks for the entire column chunk of each
+ // of the three columns "a", "b", and "c"
+ let ranges = expect_needs_data(decoder.try_decode());
+ assert_eq!(ranges, vec![4..1860, 1860..3716, 3716..11062]);
+
+ // Compute the ranges that cover only the first data page of each
+ // column (and the dictionary page which precedes it)
+ let page_index = metadata.page_index_for_row_group(0);
+ let row_group = metadata.row_group(0);
+ let mut first_page_ranges = vec![];
+ let mut second_page_ranges = vec![];
+ for (idx, column) in row_group.columns().iter().enumerate() {
+ let (start, len) = column.byte_range();
+ let locations = page_index.page_locations(idx).unwrap();
+ assert_eq!(locations.len(), 2, "expected 2 data pages per column
chunk");
+ let second_page_start = locations[1].offset as u64;
+ first_page_ranges.push(start..second_page_start);
+ second_page_ranges.push(second_page_start..start + len);
+ }
+ // Note the first range for each column includes the dictionary page
+ assert_eq!(first_page_ranges, vec![4..1734, 1860..3590, 3716..10936]);
+ assert_eq!(
+ second_page_ranges,
+ vec![1734..1860, 3590..3716, 10936..11062]
+ );
+
+ // Push only the first page of each column. This is all the data
+ // needed to decode the first batch of 100 rows.
+ push_ranges_to_decoder(&mut decoder, first_page_ranges);
+
+ // Note will likely change as part of
+ // <https://github.com/apache/arrow-rs/issues/6946>
+
+ // decoder still reports it needs the (complete) ranges
+ // it originally asked for, and does not produce a batch.
+ let ranges = expect_needs_data(decoder.try_decode());
+ assert_eq!(ranges, vec![4..1860, 1860..3716, 3716..11062]);
+
+ // Pushing the second pages as separate ranges does not help either:
+ // the decoder does not coalesce adjacent pushed ranges, so a requested
+ // range is only satisfied by a single pushed buffer that covers it.
+ push_ranges_to_decoder(&mut decoder, second_page_ranges);
+ let ranges = expect_needs_data(decoder.try_decode());
+ assert_eq!(ranges, vec![4..1860, 1860..3716, 3716..11062]);
+
+ // Only once the exact ranges originally requested are pushed does the
+ // decoder produce batches.
+ push_ranges_to_decoder(&mut decoder, ranges);
+ let batch = expect_data(decoder.try_decode());
+ assert_eq!(batch, TEST_BATCH.slice(0, 100));
+ let batch = expect_data(decoder.try_decode());
+ assert_eq!(batch, TEST_BATCH.slice(100, 100));
+ }
+
/// Decode multiple columns "a" and "b", expect that the decoder requests
/// only a single request per row group
#[test]
@@ -2706,7 +2774,7 @@ mod test {
}
/// return the metadata for the test file
- pub fn test_file_parquet_metadata() ->
Arc<crate::file::metadata::ParquetMetaData> {
+ pub fn test_file_parquet_metadata() -> Arc<ParquetMetaData> {
let mut metadata_decoder =
ParquetMetaDataPushDecoder::try_new(test_file_len()).unwrap();
push_ranges_to_metadata_decoder(&mut metadata_decoder,
vec![test_file_range()]);
let metadata = metadata_decoder.try_decode().unwrap();
@@ -2716,6 +2784,19 @@ mod test {
Arc::new(metadata)
}
+ /// return the metadata for the test file, including the offset index
+ fn test_file_parquet_metadata_with_offset_index() -> Arc<ParquetMetaData> {
+ let mut metadata_decoder =
ParquetMetaDataPushDecoder::try_new(test_file_len())
+ .unwrap()
+ .with_offset_index_policy(PageIndexPolicy::Required);
+ push_ranges_to_metadata_decoder(&mut metadata_decoder,
vec![test_file_range()]);
+ let metadata = metadata_decoder.try_decode().unwrap();
+ let DecodeResult::Data(metadata) = metadata else {
+ panic!("Expected metadata to be decoded successfully");
+ };
+ Arc::new(metadata)
+ }
+
/// Push the given ranges to the metadata decoder, simulating reading from
a file
fn push_ranges_to_metadata_decoder(
metadata_decoder: &mut ParquetMetaDataPushDecoder,