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 2f2360d97a parquet: Fix page index loading with suffix prefetch 
(#11208)
2f2360d97a is described below

commit 2f2360d97a7ec1fb3e35f33c5858fcf9289b4a09
Author: Ed Seidl <[email protected]>
AuthorDate: Fri Sep 25 07:50:17 2026 -0700

    parquet: Fix page index loading with suffix prefetch (#11208)
    
    # Which issue does this PR close?
    
    - Closes #11207.
    
    # Rationale for this change
    
    `load_metadata_via_suffix` treated bytes preceding the footer metadata
    as beginning at file offset zero, despite `MetadataSuffixFetch` not
    providing the total file size. Page-index loading could consequently
    slice the wrong bytes or fail with an out-of-bounds range error.
    
    # What changes are included in this PR?
    
    - Stop reusing suffix bytes whose absolute file offset is unknown.
    - Fetch page-index data using its actual file range.
    - Add a regression test covering page-index loading through
    `load_via_suffix_and_finish` with a large prefetch hint.
    
    # Are these changes tested?
    
    Yes. The new regression test fails before the fix with:
    
    ```text
    Corrupted parquet file: index data range (...) exceeds remainder length 
(...)
    ```
    
    and passes afterward.
    
    The focused async metadata-reader suite passes:
    
    ```text
    cargo test -p parquet --lib 'file::metadata::reader::async_tests' 
--features arrow,async
    ```
    
    # Are there any user-facing changes?
    
    Page-index loading through `MetadataSuffixFetch` now performs a separate
    range fetch when the prefetched bytes have no known absolute offset.
    This prevents incorrect page-index decoding. There are no public API
    changes.
    
    # AI assistance
    
    I used OpenAI Codex to help investigate the failure, draft the
    regression test, and implement the fix. I reviewed and understand the
    resulting changes and take responsibility for them.
---
 parquet/src/file/metadata/reader.rs | 46 ++++++++++++++++++++++++++++++-------
 1 file changed, 38 insertions(+), 8 deletions(-)

diff --git a/parquet/src/file/metadata/reader.rs 
b/parquet/src/file/metadata/reader.rs
index e4cec0fd0d..3ffce28523 100644
--- a/parquet/src/file/metadata/reader.rs
+++ b/parquet/src/file/metadata/reader.rs
@@ -452,7 +452,7 @@ impl ParquetMetaDataReader {
         &mut self,
         mut fetch: F,
     ) -> Result<()> {
-        let (metadata, remainder) = self.load_metadata_via_suffix(&mut 
fetch).await?;
+        let metadata = self.load_metadata_via_suffix(&mut fetch).await?;
 
         self.metadata = Some(metadata);
 
@@ -462,7 +462,7 @@ impl ParquetMetaDataReader {
             return Ok(());
         }
 
-        self.load_page_index_with_remainder(fetch, remainder).await
+        self.load_page_index_with_remainder(fetch, None).await
     }
 
     /// Asynchronously fetch the page index structures when a 
[`ParquetMetaData`] has already
@@ -643,10 +643,13 @@ impl ParquetMetaDataReader {
     }
 
     #[cfg(all(feature = "async", feature = "arrow"))]
+    // Unlike load_metadata, the file size is not known so it is not safe
+    // to use any leftover bytes that may have been pre-fetched. Thus this
+    // returns only the metadata.
     async fn load_metadata_via_suffix<F: MetadataSuffixFetch>(
         &self,
         fetch: &mut F,
-    ) -> Result<(ParquetMetaData, Option<(usize, Bytes)>)> {
+    ) -> Result<ParquetMetaData> {
         let prefetch = self.get_prefetch_size();
 
         let suffix = fetch.fetch_suffix(prefetch).await?;
@@ -684,14 +687,11 @@ impl ParquetMetaDataReader {
 
             // need to slice off the footer or decryption fails
             let meta = meta.slice(0..length);
-            Ok((self.decode_footer_metadata(meta, file_size, footer)?, None))
+            Ok(self.decode_footer_metadata(meta, file_size, footer)?)
         } else {
             let metadata_start = suffix_len - metadata_offset;
             let slice = suffix.slice(metadata_start..suffix_len - FOOTER_SIZE);
-            Ok((
-                self.decode_footer_metadata(slice, file_size, footer)?,
-                Some((0, suffix.slice(..metadata_start))),
-            ))
+            Ok(self.decode_footer_metadata(slice, file_size, footer)?)
         }
     }
 
@@ -1347,6 +1347,36 @@ mod async_tests {
         assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
     }
 
+    #[tokio::test]
+    async fn test_page_index_via_suffix_with_prefetch() {
+        let mut file = get_test_file("alltypes_tiny_pages.parquet");
+        let mut suffix_file = file.try_clone().unwrap();
+        let len = file.len();
+        let fetch_count = AtomicUsize::new(0);
+        let suffix_fetch_count = AtomicUsize::new(0);
+
+        let mut fetch = |range| {
+            fetch_count.fetch_add(1, Ordering::SeqCst);
+            futures::future::ready(read_range(&mut file, range))
+        };
+        let mut suffix_fetch = |suffix| {
+            suffix_fetch_count.fetch_add(1, Ordering::SeqCst);
+            futures::future::ready(read_suffix(&mut suffix_file, suffix))
+        };
+
+        let input = MetadataSuffixFetchFn(&mut fetch, &mut suffix_fetch);
+        let metadata = ParquetMetaDataReader::new()
+            .with_page_index_policy(PageIndexPolicy::Required)
+            .with_prefetch_hint(Some((len - 1000) as usize))
+            .load_via_suffix_and_finish(input)
+            .await
+            .unwrap();
+
+        assert_eq!(suffix_fetch_count.load(Ordering::SeqCst), 1);
+        assert_eq!(fetch_count.load(Ordering::SeqCst), 1);
+        assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
+    }
+
     fn write_parquet_file(offset_index_disabled: bool) -> 
Result<NamedTempFile> {
         let schema = Arc::new(Schema::new(vec![Field::new("a", 
DataType::Int32, false)]));
         let batch = RecordBatch::try_new(

Reply via email to