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(