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,

Reply via email to