adriangb commented on code in PR #11132: URL: https://github.com/apache/arrow-rs/pull/11132#discussion_r4077669477
########## parquet/tests/arrow_reader/custom_page_index_provider.rs: ########## @@ -0,0 +1,460 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! End-to-end test for custom PageIndexProvider implementation. +//! +//! This test validates reading with a custom PageIndexProvider through the +//! actual read path, ensuring that: +//! - Columns with page indexes use offset-index-driven fetch +//! - Columns without page indexes fall back to whole-column-chunk fetching +//! - Correct results are returned in both cases + +use arrow_array::{Array, Int32Array, RecordBatch, StringArray}; +use arrow_schema::{DataType, Field, Schema}; +use bytes::Bytes; +use parquet::arrow::ArrowWriter; +use parquet::arrow::arrow_reader::{ + ArrowReaderMetadata, ArrowReaderOptions, ParquetRecordBatchReaderBuilder, RowSelection, + RowSelectionPolicy, +}; +use parquet::file::metadata::page_index::PageIndexProvider; +use parquet::file::metadata::{PageIndexPolicy, ParquetMetaData}; +use parquet::file::page_index::column_index::ColumnIndexMetaData; +use parquet::file::page_index::index_reader::{decode_column_index, decode_offset_index}; +use parquet::file::page_index::offset_index::OffsetIndexMetaData; +use parquet::file::properties::{EnabledStatistics, WriterProperties}; +use std::collections::HashMap; +use std::collections::hash_map::Entry; +use std::fs::File; +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use tempfile::NamedTempFile; + +#[test] +fn test_read_with_custom_page_index_provider() { + // Step 1: Write a parquet file with page indexes + let temp_file = create_test_file(); + let file_bytes = Bytes::from(std::fs::read(temp_file.path()).unwrap()); + + // Step 2: Load metadata WITHOUT page indexes initially + let file = File::open(temp_file.path()).unwrap(); + let builder = ParquetRecordBatchReaderBuilder::try_new_with_options( + file, + ArrowReaderOptions::default().with_page_index_policy(PageIndexPolicy::Skip), + ) + .unwrap(); + + let metadata = builder.metadata().as_ref(); + + // Verify we have multiple row groups + assert_eq!(metadata.num_row_groups(), 3); + + // Step 3: Create a custom provider and populate it selectively + // Simulate a scenario where: + // - We only populate indexes for row groups 0 and 2 (skipping row group 1) + // - For row group 0: populate column 0 (id) and column 1 (value) + // - For row group 2: populate column 0 (id) only + let mut provider = SelectivePageIndexProvider::new(file_bytes); + + // Populate indexes for row group 0, columns 0 and 1 + provider.fetch_column_index(0, 0, metadata).unwrap(); + provider.fetch_offset_index(0, 0, metadata).unwrap(); + provider.fetch_column_index(0, 1, metadata).unwrap(); + provider.fetch_offset_index(0, 1, metadata).unwrap(); + + // Populate indexes for row group 2, column 0 only + provider.fetch_column_index(2, 0, metadata).unwrap(); + provider.fetch_offset_index(2, 0, metadata).unwrap(); + + // Verify the provider has the expected indexes + assert!(provider.column_index(0, 0).is_some()); + assert!(provider.offset_index(0, 0).is_some()); + assert!(provider.column_index(0, 1).is_some()); + assert!(provider.offset_index(0, 1).is_some()); + assert!(provider.column_index(0, 2).is_none()); // Not populated + assert!(provider.column_index(2, 0).is_some()); + assert!(provider.column_index(2, 1).is_none()); // Not populated + + // Reset statistics so validation checks don't affect the final counts + provider.reset_stats(); + + // Step 4: Install the custom provider into metadata + let mut metadata_builder = metadata.clone().into_builder(); + metadata_builder = metadata_builder.set_page_index(Some(Arc::new(provider))); + let metadata_with_custom_index = Arc::new(metadata_builder.build()); + + // Step 5: Create ArrowReaderMetadata with the custom page index + let arrow_metadata = ArrowReaderMetadata::try_new( + metadata_with_custom_index.clone(), + ArrowReaderOptions::default(), + ) + .unwrap(); + + // Step 6: Read data with RowSelection that triggers page skipping + let file = File::open(temp_file.path()).unwrap(); + let builder = ParquetRecordBatchReaderBuilder::new_with_metadata(file, arrow_metadata); + + // Create a RowSelection that: + // - Selects rows 20-30 (in row group 0) + // - Selects rows 60-70 (in row group 1) + // - Selects rows 120-130 (in row group 2) + let selection = RowSelection::from(vec![ + // Skip first 20 rows + parquet::arrow::arrow_reader::RowSelector::skip(20), + // Select rows 20-30 + parquet::arrow::arrow_reader::RowSelector::select(10), + // Skip rows 30-60 + parquet::arrow::arrow_reader::RowSelector::skip(30), + // Select rows 60-70 + parquet::arrow::arrow_reader::RowSelector::select(10), + // Skip rows 70-120 + parquet::arrow::arrow_reader::RowSelector::skip(50), + // Select rows 120-130 + parquet::arrow::arrow_reader::RowSelector::select(10), + ]); + + let reader = builder + .with_row_selection(selection) + // make sure we're actually using the page index and not using a mask + .with_row_selection_policy(RowSelectionPolicy::Selectors) + .build() + .unwrap(); + + // Collect all batches + let batches: Vec<RecordBatch> = reader.collect::<Result<Vec<_>, _>>().unwrap(); + + // Step 7: Check statistics + // Note: We need to get the provider reference from the metadata to check stats + let provider_ref = metadata_with_custom_index + .page_index() + .unwrap() + .as_any() + .downcast_ref::<SelectivePageIndexProvider>() + .expect("Expected SelectivePageIndexProvider"); + + // Check that the custom provider was used + let (hits, misses) = provider_ref.stats(); + assert!(hits > 0, "provider should have hits"); + assert!(misses > 0, "provider should have misses"); + + // Since we didn't use predicates, the column index should be untouched + let (col_hits, col_misses) = provider_ref.column_index_stats(); + assert_eq!(col_hits, 0, "column index should not be used"); + assert_eq!(col_misses, 0, "column index should not be used"); + + let (off_hits, off_misses) = provider_ref.offset_index_stats(); + assert!(off_hits > 0, "offset index should have hits"); + assert!(off_misses > 0, "offset index should have misses"); Review Comment: **Test gap (B3).** These asserts show that the reader called `offset_index()`. They do not show that the reader used the result. Check: I changed `ReaderPageIterator::next_page_reader` (`arrow_reader/mod.rs` L1466) to pass `None` instead of `page_locations`. This test still passes. So nothing checks the two claims in the module doc (L22-23). The push decoder can check these claims directly, with the byte ranges that it requests. With the fix in [#11182](https://github.com/apache/arrow-rs/pull/11182), each chunk with an offset index fetches 245 of 363 bytes. Each chunk without an offset index fetches the whole chunk. ```rust // in the push loop: requested.extend(ranges.iter().cloned()); let fetched = |rg: usize, col: usize| -> u64 { let (start, len) = metadata.row_group(rg).column(col).byte_range(); requested .iter() .filter(|r| r.start >= start && r.end <= start + len) .map(|r| r.end - r.start) .sum() }; let chunk_len = |rg: usize, col: usize| metadata.row_group(rg).column(col).byte_range().1; assert!(fetched(0, 0) < chunk_len(0, 0)); // offset index: selected pages only assert_eq!(fetched(0, 2), chunk_len(0, 2)); // no offset index: whole chunk ``` If you keep the counters: L150-152 (`stats()`) repeat L159-161, so you can remove `stats()`. ########## parquet/tests/arrow_reader/custom_page_index_provider.rs: ########## @@ -0,0 +1,460 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! End-to-end test for custom PageIndexProvider implementation. +//! +//! This test validates reading with a custom PageIndexProvider through the +//! actual read path, ensuring that: +//! - Columns with page indexes use offset-index-driven fetch +//! - Columns without page indexes fall back to whole-column-chunk fetching +//! - Correct results are returned in both cases + +use arrow_array::{Array, Int32Array, RecordBatch, StringArray}; +use arrow_schema::{DataType, Field, Schema}; +use bytes::Bytes; +use parquet::arrow::ArrowWriter; +use parquet::arrow::arrow_reader::{ + ArrowReaderMetadata, ArrowReaderOptions, ParquetRecordBatchReaderBuilder, RowSelection, + RowSelectionPolicy, +}; +use parquet::file::metadata::page_index::PageIndexProvider; +use parquet::file::metadata::{PageIndexPolicy, ParquetMetaData}; +use parquet::file::page_index::column_index::ColumnIndexMetaData; +use parquet::file::page_index::index_reader::{decode_column_index, decode_offset_index}; +use parquet::file::page_index::offset_index::OffsetIndexMetaData; +use parquet::file::properties::{EnabledStatistics, WriterProperties}; +use std::collections::HashMap; +use std::collections::hash_map::Entry; +use std::fs::File; +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use tempfile::NamedTempFile; + +#[test] +fn test_read_with_custom_page_index_provider() { + // Step 1: Write a parquet file with page indexes + let temp_file = create_test_file(); + let file_bytes = Bytes::from(std::fs::read(temp_file.path()).unwrap()); + + // Step 2: Load metadata WITHOUT page indexes initially + let file = File::open(temp_file.path()).unwrap(); + let builder = ParquetRecordBatchReaderBuilder::try_new_with_options( + file, + ArrowReaderOptions::default().with_page_index_policy(PageIndexPolicy::Skip), + ) + .unwrap(); + + let metadata = builder.metadata().as_ref(); + + // Verify we have multiple row groups + assert_eq!(metadata.num_row_groups(), 3); + + // Step 3: Create a custom provider and populate it selectively + // Simulate a scenario where: + // - We only populate indexes for row groups 0 and 2 (skipping row group 1) + // - For row group 0: populate column 0 (id) and column 1 (value) + // - For row group 2: populate column 0 (id) only + let mut provider = SelectivePageIndexProvider::new(file_bytes); + + // Populate indexes for row group 0, columns 0 and 1 + provider.fetch_column_index(0, 0, metadata).unwrap(); + provider.fetch_offset_index(0, 0, metadata).unwrap(); + provider.fetch_column_index(0, 1, metadata).unwrap(); + provider.fetch_offset_index(0, 1, metadata).unwrap(); + + // Populate indexes for row group 2, column 0 only + provider.fetch_column_index(2, 0, metadata).unwrap(); + provider.fetch_offset_index(2, 0, metadata).unwrap(); Review Comment: **Test gap (B2).** In each row group, the chunks without an index come after the chunks with an index (rg 0: cols 2, 3; rg 1: all; rg 2: cols 1-3). The bug in the review body is positional, so this layout does not show if a fix is complete. Example: a naive fix ("when no offsets are left, use the whole chunk") makes this test pass with all three readers. With a gap before an indexed column, the same naive fix fails in async and push: ```text ArrowError("Parquet argument error: Parquet error: Invalid offset in sparse column chunk data: 3861, no matching page found. [...]") ``` This change adds such a gap in row group 2. The asserts at L89-90 stay true. Results with this change: sync passes on this head. Async and push pass with the correct fix, and fail with the naive fix. ```suggestion // - For row group 2: populate columns 0 (id) and 2 (name), so column 1 is a gap let mut provider = SelectivePageIndexProvider::new(file_bytes); // Populate indexes for row group 0, columns 0 and 1 provider.fetch_column_index(0, 0, metadata).unwrap(); provider.fetch_offset_index(0, 0, metadata).unwrap(); provider.fetch_column_index(0, 1, metadata).unwrap(); provider.fetch_offset_index(0, 1, metadata).unwrap(); // Populate indexes for row group 2, columns 0 and 2 (column 1 has no index) provider.fetch_column_index(2, 0, metadata).unwrap(); provider.fetch_offset_index(2, 0, metadata).unwrap(); provider.fetch_column_index(2, 2, metadata).unwrap(); provider.fetch_offset_index(2, 2, metadata).unwrap(); ``` ########## parquet/tests/arrow_reader/custom_page_index_provider.rs: ########## @@ -0,0 +1,460 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! End-to-end test for custom PageIndexProvider implementation. +//! +//! This test validates reading with a custom PageIndexProvider through the +//! actual read path, ensuring that: +//! - Columns with page indexes use offset-index-driven fetch +//! - Columns without page indexes fall back to whole-column-chunk fetching +//! - Correct results are returned in both cases + +use arrow_array::{Array, Int32Array, RecordBatch, StringArray}; +use arrow_schema::{DataType, Field, Schema}; +use bytes::Bytes; +use parquet::arrow::ArrowWriter; +use parquet::arrow::arrow_reader::{ + ArrowReaderMetadata, ArrowReaderOptions, ParquetRecordBatchReaderBuilder, RowSelection, + RowSelectionPolicy, +}; +use parquet::file::metadata::page_index::PageIndexProvider; +use parquet::file::metadata::{PageIndexPolicy, ParquetMetaData}; +use parquet::file::page_index::column_index::ColumnIndexMetaData; +use parquet::file::page_index::index_reader::{decode_column_index, decode_offset_index}; +use parquet::file::page_index::offset_index::OffsetIndexMetaData; +use parquet::file::properties::{EnabledStatistics, WriterProperties}; +use std::collections::HashMap; +use std::collections::hash_map::Entry; +use std::fs::File; +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use tempfile::NamedTempFile; + +#[test] +fn test_read_with_custom_page_index_provider() { + // Step 1: Write a parquet file with page indexes + let temp_file = create_test_file(); + let file_bytes = Bytes::from(std::fs::read(temp_file.path()).unwrap()); + + // Step 2: Load metadata WITHOUT page indexes initially + let file = File::open(temp_file.path()).unwrap(); + let builder = ParquetRecordBatchReaderBuilder::try_new_with_options( + file, + ArrowReaderOptions::default().with_page_index_policy(PageIndexPolicy::Skip), + ) + .unwrap(); + + let metadata = builder.metadata().as_ref(); + + // Verify we have multiple row groups + assert_eq!(metadata.num_row_groups(), 3); + + // Step 3: Create a custom provider and populate it selectively + // Simulate a scenario where: + // - We only populate indexes for row groups 0 and 2 (skipping row group 1) + // - For row group 0: populate column 0 (id) and column 1 (value) + // - For row group 2: populate column 0 (id) only + let mut provider = SelectivePageIndexProvider::new(file_bytes); + + // Populate indexes for row group 0, columns 0 and 1 + provider.fetch_column_index(0, 0, metadata).unwrap(); + provider.fetch_offset_index(0, 0, metadata).unwrap(); + provider.fetch_column_index(0, 1, metadata).unwrap(); + provider.fetch_offset_index(0, 1, metadata).unwrap(); + + // Populate indexes for row group 2, column 0 only + provider.fetch_column_index(2, 0, metadata).unwrap(); + provider.fetch_offset_index(2, 0, metadata).unwrap(); + + // Verify the provider has the expected indexes + assert!(provider.column_index(0, 0).is_some()); + assert!(provider.offset_index(0, 0).is_some()); + assert!(provider.column_index(0, 1).is_some()); + assert!(provider.offset_index(0, 1).is_some()); + assert!(provider.column_index(0, 2).is_none()); // Not populated + assert!(provider.column_index(2, 0).is_some()); + assert!(provider.column_index(2, 1).is_none()); // Not populated + + // Reset statistics so validation checks don't affect the final counts + provider.reset_stats(); + + // Step 4: Install the custom provider into metadata + let mut metadata_builder = metadata.clone().into_builder(); + metadata_builder = metadata_builder.set_page_index(Some(Arc::new(provider))); + let metadata_with_custom_index = Arc::new(metadata_builder.build()); + + // Step 5: Create ArrowReaderMetadata with the custom page index + let arrow_metadata = ArrowReaderMetadata::try_new( + metadata_with_custom_index.clone(), + ArrowReaderOptions::default(), + ) + .unwrap(); + + // Step 6: Read data with RowSelection that triggers page skipping + let file = File::open(temp_file.path()).unwrap(); + let builder = ParquetRecordBatchReaderBuilder::new_with_metadata(file, arrow_metadata); + + // Create a RowSelection that: + // - Selects rows 20-30 (in row group 0) + // - Selects rows 60-70 (in row group 1) + // - Selects rows 120-130 (in row group 2) + let selection = RowSelection::from(vec![ + // Skip first 20 rows + parquet::arrow::arrow_reader::RowSelector::skip(20), + // Select rows 20-30 + parquet::arrow::arrow_reader::RowSelector::select(10), + // Skip rows 30-60 + parquet::arrow::arrow_reader::RowSelector::skip(30), + // Select rows 60-70 + parquet::arrow::arrow_reader::RowSelector::select(10), + // Skip rows 70-120 + parquet::arrow::arrow_reader::RowSelector::skip(50), + // Select rows 120-130 + parquet::arrow::arrow_reader::RowSelector::select(10), + ]); + + let reader = builder + .with_row_selection(selection) + // make sure we're actually using the page index and not using a mask + .with_row_selection_policy(RowSelectionPolicy::Selectors) + .build() + .unwrap(); + + // Collect all batches + let batches: Vec<RecordBatch> = reader.collect::<Result<Vec<_>, _>>().unwrap(); Review Comment: **Bug (pre-existing) / Test gap (B1).** This reads only with the sync reader. The same `arrow_metadata` and `selection` fail with the other two readers on this head: ```text async: Err(General("Invalid column index 2, column was not fetched")) push: Err(General("Invalid column index 2, column was not fetched")) ``` Cause: `InMemoryRowGroup` aligns `page_start_offsets` to the columns by position (see the review body). The sync reader does not use `InMemoryRowGroup`. Also, the module doc (L22-23) is about "fetch", and only the async reader and the push decoder fetch ranges. Suggestion (optional for this PR): run the same test body with each reader. With the fix, all three pass. Without it, only `Sync` passes. If this PR merges before the fix, add the push and async variants with `#[ignore = "https://github.com/apache/arrow-rs/pull/11182"]`, or add them in the fix PR ([#11182](https://github.com/apache/arrow-rs/pull/11182)). <details><summary>Snippet: one test for each reader (compiled and run on this head)</summary> ```rust use parquet::DecodeResult; use parquet::arrow::push_decoder::ParquetPushDecoderBuilder; #[derive(Debug, Clone, Copy)] enum Reader { Sync, Push, #[cfg(feature = "async")] Async, } /// Read `file` with `reader`, using `metadata` and `selection` fn read( reader: Reader, file: Bytes, metadata: ArrowReaderMetadata, selection: RowSelection, ) -> Vec<RecordBatch> { match reader { Reader::Sync => ParquetRecordBatchReaderBuilder::new_with_metadata(file, metadata) .with_row_selection(selection) .with_row_selection_policy(RowSelectionPolicy::Selectors) .build() .unwrap() .collect::<Result<_, _>>() .unwrap(), Reader::Push => { let mut decoder = ParquetPushDecoderBuilder::new_with_metadata(metadata) .with_row_selection(selection) .with_row_selection_policy(RowSelectionPolicy::Selectors) .build() .unwrap(); let mut batches = vec![]; loop { match decoder.try_decode().unwrap() { DecodeResult::NeedsData(ranges) => { let data = ranges .iter() .map(|r| file.slice(r.start as usize..r.end as usize)) .collect(); decoder.push_ranges(ranges, data).unwrap(); } DecodeResult::Data(batch) => batches.push(batch), DecodeResult::Finished => return batches, } } } #[cfg(feature = "async")] Reader::Async => { use futures::TryStreamExt; let input = std::io::Cursor::new(file.to_vec()); let rt = tokio::runtime::Builder::new_current_thread() .build() .unwrap(); rt.block_on(async { parquet::arrow::ParquetRecordBatchStreamBuilder::new_with_metadata(input, metadata) .with_row_selection(selection) .with_row_selection_policy(RowSelectionPolicy::Selectors) .build() .unwrap() .try_collect() .await .unwrap() }) } } } #[test] fn test_read_with_custom_page_index_provider_sync() { run_test(Reader::Sync); } #[test] fn test_read_with_custom_page_index_provider_push() { run_test(Reader::Push); } #[cfg(feature = "async")] #[test] fn test_read_with_custom_page_index_provider_async() { run_test(Reader::Async); } ``` In `fn run_test(reader: Reader)` (the current test body), remove L108-109, and L130-138 become one line: ```rust let batches = read(reader, file_bytes.clone(), arrow_metadata, selection); ``` L71 then needs `SelectivePageIndexProvider::new(file_bytes.clone())`. The push variant has no feature gate, so the default CI job runs it too. </details> ########## parquet/tests/arrow_reader/custom_page_index_provider.rs: ########## @@ -0,0 +1,460 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! End-to-end test for custom PageIndexProvider implementation. +//! +//! This test validates reading with a custom PageIndexProvider through the +//! actual read path, ensuring that: +//! - Columns with page indexes use offset-index-driven fetch +//! - Columns without page indexes fall back to whole-column-chunk fetching +//! - Correct results are returned in both cases + +use arrow_array::{Array, Int32Array, RecordBatch, StringArray}; +use arrow_schema::{DataType, Field, Schema}; +use bytes::Bytes; +use parquet::arrow::ArrowWriter; +use parquet::arrow::arrow_reader::{ + ArrowReaderMetadata, ArrowReaderOptions, ParquetRecordBatchReaderBuilder, RowSelection, + RowSelectionPolicy, +}; +use parquet::file::metadata::page_index::PageIndexProvider; +use parquet::file::metadata::{PageIndexPolicy, ParquetMetaData}; +use parquet::file::page_index::column_index::ColumnIndexMetaData; +use parquet::file::page_index::index_reader::{decode_column_index, decode_offset_index}; +use parquet::file::page_index::offset_index::OffsetIndexMetaData; +use parquet::file::properties::{EnabledStatistics, WriterProperties}; +use std::collections::HashMap; +use std::collections::hash_map::Entry; +use std::fs::File; +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use tempfile::NamedTempFile; + +#[test] +fn test_read_with_custom_page_index_provider() { + // Step 1: Write a parquet file with page indexes + let temp_file = create_test_file(); + let file_bytes = Bytes::from(std::fs::read(temp_file.path()).unwrap()); + + // Step 2: Load metadata WITHOUT page indexes initially + let file = File::open(temp_file.path()).unwrap(); + let builder = ParquetRecordBatchReaderBuilder::try_new_with_options( + file, + ArrowReaderOptions::default().with_page_index_policy(PageIndexPolicy::Skip), + ) + .unwrap(); + + let metadata = builder.metadata().as_ref(); + + // Verify we have multiple row groups + assert_eq!(metadata.num_row_groups(), 3); + + // Step 3: Create a custom provider and populate it selectively + // Simulate a scenario where: + // - We only populate indexes for row groups 0 and 2 (skipping row group 1) + // - For row group 0: populate column 0 (id) and column 1 (value) + // - For row group 2: populate column 0 (id) only + let mut provider = SelectivePageIndexProvider::new(file_bytes); + + // Populate indexes for row group 0, columns 0 and 1 + provider.fetch_column_index(0, 0, metadata).unwrap(); + provider.fetch_offset_index(0, 0, metadata).unwrap(); + provider.fetch_column_index(0, 1, metadata).unwrap(); + provider.fetch_offset_index(0, 1, metadata).unwrap(); + + // Populate indexes for row group 2, column 0 only + provider.fetch_column_index(2, 0, metadata).unwrap(); + provider.fetch_offset_index(2, 0, metadata).unwrap(); + + // Verify the provider has the expected indexes + assert!(provider.column_index(0, 0).is_some()); + assert!(provider.offset_index(0, 0).is_some()); + assert!(provider.column_index(0, 1).is_some()); + assert!(provider.offset_index(0, 1).is_some()); + assert!(provider.column_index(0, 2).is_none()); // Not populated + assert!(provider.column_index(2, 0).is_some()); + assert!(provider.column_index(2, 1).is_none()); // Not populated + + // Reset statistics so validation checks don't affect the final counts + provider.reset_stats(); + + // Step 4: Install the custom provider into metadata + let mut metadata_builder = metadata.clone().into_builder(); + metadata_builder = metadata_builder.set_page_index(Some(Arc::new(provider))); + let metadata_with_custom_index = Arc::new(metadata_builder.build()); + + // Step 5: Create ArrowReaderMetadata with the custom page index + let arrow_metadata = ArrowReaderMetadata::try_new( + metadata_with_custom_index.clone(), + ArrowReaderOptions::default(), + ) + .unwrap(); + + // Step 6: Read data with RowSelection that triggers page skipping + let file = File::open(temp_file.path()).unwrap(); + let builder = ParquetRecordBatchReaderBuilder::new_with_metadata(file, arrow_metadata); + + // Create a RowSelection that: + // - Selects rows 20-30 (in row group 0) + // - Selects rows 60-70 (in row group 1) + // - Selects rows 120-130 (in row group 2) + let selection = RowSelection::from(vec![ + // Skip first 20 rows + parquet::arrow::arrow_reader::RowSelector::skip(20), + // Select rows 20-30 + parquet::arrow::arrow_reader::RowSelector::select(10), + // Skip rows 30-60 + parquet::arrow::arrow_reader::RowSelector::skip(30), + // Select rows 60-70 + parquet::arrow::arrow_reader::RowSelector::select(10), + // Skip rows 70-120 + parquet::arrow::arrow_reader::RowSelector::skip(50), + // Select rows 120-130 + parquet::arrow::arrow_reader::RowSelector::select(10), + ]); + + let reader = builder + .with_row_selection(selection) + // make sure we're actually using the page index and not using a mask + .with_row_selection_policy(RowSelectionPolicy::Selectors) + .build() + .unwrap(); + + // Collect all batches + let batches: Vec<RecordBatch> = reader.collect::<Result<Vec<_>, _>>().unwrap(); + + // Step 7: Check statistics + // Note: We need to get the provider reference from the metadata to check stats + let provider_ref = metadata_with_custom_index + .page_index() + .unwrap() + .as_any() + .downcast_ref::<SelectivePageIndexProvider>() + .expect("Expected SelectivePageIndexProvider"); + + // Check that the custom provider was used + let (hits, misses) = provider_ref.stats(); + assert!(hits > 0, "provider should have hits"); + assert!(misses > 0, "provider should have misses"); + + // Since we didn't use predicates, the column index should be untouched + let (col_hits, col_misses) = provider_ref.column_index_stats(); + assert_eq!(col_hits, 0, "column index should not be used"); + assert_eq!(col_misses, 0, "column index should not be used"); + + let (off_hits, off_misses) = provider_ref.offset_index_stats(); + assert!(off_hits > 0, "offset index should have hits"); + assert!(off_misses > 0, "offset index should have misses"); + + // Step 8: Verify correct results + // We expect 30 rows total: 10 from each selected range + let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum(); + assert_eq!(total_rows, 30); + + // Verify the actual data values + let mut all_ids = Vec::new(); + let mut all_values = Vec::new(); + let mut all_names = Vec::new(); + let mut all_scores = Vec::new(); + + for batch in batches { + let ids = batch + .column(0) + .as_any() + .downcast_ref::<Int32Array>() + .unwrap(); + let values = batch + .column(1) + .as_any() + .downcast_ref::<Int32Array>() + .unwrap(); + let names = batch + .column(2) + .as_any() + .downcast_ref::<StringArray>() + .unwrap(); + let scores = batch + .column(3) + .as_any() + .downcast_ref::<Int32Array>() + .unwrap(); + + all_ids.extend(ids.iter().flatten()); + all_values.extend(values.iter().flatten()); + all_names.extend(names.iter().flatten().map(|s| s.to_owned())); + all_scores.extend(scores.iter().flatten()); + } + + // Expected values for the selected rows: + // rows 20-30: id=20..30, value=40..60, score=60..90 + // rows 60-70: id=60..70, value=120..140, score=180..210 + // rows 120-130: id=120..130, value=240..260, score=360..390 + let expected_ids: Vec<i32> = (20..30).chain(60..70).chain(120..130).collect(); + let expected_values: Vec<i32> = expected_ids.iter().map(|x| x * 2).collect(); + let expected_scores: Vec<i32> = expected_ids.iter().map(|x| x * 3).collect(); + let expected_names: Vec<String> = expected_ids.iter().map(|x| format!("name{x}")).collect(); + + assert_eq!(all_ids, expected_ids, "IDs don't match expected values"); + assert_eq!( + all_values, expected_values, + "Values don't match expected values" + ); + assert_eq!( + all_scores, expected_scores, + "Scores don't match expected values" + ); + assert_eq!( + all_names, expected_names, + "Names don't match expected values" + ); +} + +/// A custom PageIndexProvider that only stores indexes for a subset of columns +/// +/// This provider mimics a real-world scenario where an application has loaded +/// page indexes selectively based on query predicates and projections. +#[derive(Debug)] +struct SelectivePageIndexProvider { + file_bytes: Bytes, + // indexes are accessed first by row_group index and then by column index + column_indexes: Option<HashMap<usize, HashMap<usize, ColumnIndexMetaData>>>, + offset_indexes: Option<HashMap<usize, HashMap<usize, OffsetIndexMetaData>>>, + // Usage statistics + column_index_hits: Arc<AtomicUsize>, + column_index_misses: Arc<AtomicUsize>, + offset_index_hits: Arc<AtomicUsize>, + offset_index_misses: Arc<AtomicUsize>, +} Review Comment: **Nit (B4).** The helper is about 160 of the 460 lines. It can be much smaller: - Use one key: `HashMap<(usize, usize), T>` instead of `Option<HashMap<usize, HashMap<usize, T>>>`. - Use plain `AtomicUsize` fields. Nothing clones the `Arc`s. - `pub fn` on a private test type has no effect. - L83-90 test the helper, not the reader. Without them, `reset_stats()` is not necessary. - Build the provider from a list. Then the `file_bytes` field and the `Entry` code are not necessary: ```rust fn new(file: &Bytes, metadata: &ParquetMetaData, chunks: &[(usize, usize)]) -> Self ``` My scratch version of the same provider (struct, trait impl and constructor, no counters) has about 40 lines. ########## parquet/tests/arrow_reader/custom_page_index_provider.rs: ########## @@ -0,0 +1,460 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! End-to-end test for custom PageIndexProvider implementation. +//! +//! This test validates reading with a custom PageIndexProvider through the +//! actual read path, ensuring that: +//! - Columns with page indexes use offset-index-driven fetch +//! - Columns without page indexes fall back to whole-column-chunk fetching +//! - Correct results are returned in both cases + +use arrow_array::{Array, Int32Array, RecordBatch, StringArray}; +use arrow_schema::{DataType, Field, Schema}; +use bytes::Bytes; +use parquet::arrow::ArrowWriter; +use parquet::arrow::arrow_reader::{ + ArrowReaderMetadata, ArrowReaderOptions, ParquetRecordBatchReaderBuilder, RowSelection, + RowSelectionPolicy, +}; +use parquet::file::metadata::page_index::PageIndexProvider; +use parquet::file::metadata::{PageIndexPolicy, ParquetMetaData}; +use parquet::file::page_index::column_index::ColumnIndexMetaData; +use parquet::file::page_index::index_reader::{decode_column_index, decode_offset_index}; +use parquet::file::page_index::offset_index::OffsetIndexMetaData; +use parquet::file::properties::{EnabledStatistics, WriterProperties}; +use std::collections::HashMap; +use std::collections::hash_map::Entry; +use std::fs::File; +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use tempfile::NamedTempFile; + +#[test] +fn test_read_with_custom_page_index_provider() { + // Step 1: Write a parquet file with page indexes + let temp_file = create_test_file(); + let file_bytes = Bytes::from(std::fs::read(temp_file.path()).unwrap()); + + // Step 2: Load metadata WITHOUT page indexes initially + let file = File::open(temp_file.path()).unwrap(); + let builder = ParquetRecordBatchReaderBuilder::try_new_with_options( + file, + ArrowReaderOptions::default().with_page_index_policy(PageIndexPolicy::Skip), + ) + .unwrap(); + + let metadata = builder.metadata().as_ref(); + + // Verify we have multiple row groups + assert_eq!(metadata.num_row_groups(), 3); + + // Step 3: Create a custom provider and populate it selectively + // Simulate a scenario where: + // - We only populate indexes for row groups 0 and 2 (skipping row group 1) + // - For row group 0: populate column 0 (id) and column 1 (value) + // - For row group 2: populate column 0 (id) only + let mut provider = SelectivePageIndexProvider::new(file_bytes); + + // Populate indexes for row group 0, columns 0 and 1 + provider.fetch_column_index(0, 0, metadata).unwrap(); + provider.fetch_offset_index(0, 0, metadata).unwrap(); + provider.fetch_column_index(0, 1, metadata).unwrap(); + provider.fetch_offset_index(0, 1, metadata).unwrap(); + + // Populate indexes for row group 2, column 0 only + provider.fetch_column_index(2, 0, metadata).unwrap(); + provider.fetch_offset_index(2, 0, metadata).unwrap(); + + // Verify the provider has the expected indexes + assert!(provider.column_index(0, 0).is_some()); + assert!(provider.offset_index(0, 0).is_some()); + assert!(provider.column_index(0, 1).is_some()); + assert!(provider.offset_index(0, 1).is_some()); + assert!(provider.column_index(0, 2).is_none()); // Not populated + assert!(provider.column_index(2, 0).is_some()); + assert!(provider.column_index(2, 1).is_none()); // Not populated + + // Reset statistics so validation checks don't affect the final counts + provider.reset_stats(); + + // Step 4: Install the custom provider into metadata + let mut metadata_builder = metadata.clone().into_builder(); + metadata_builder = metadata_builder.set_page_index(Some(Arc::new(provider))); + let metadata_with_custom_index = Arc::new(metadata_builder.build()); + + // Step 5: Create ArrowReaderMetadata with the custom page index + let arrow_metadata = ArrowReaderMetadata::try_new( + metadata_with_custom_index.clone(), + ArrowReaderOptions::default(), + ) + .unwrap(); + + // Step 6: Read data with RowSelection that triggers page skipping + let file = File::open(temp_file.path()).unwrap(); + let builder = ParquetRecordBatchReaderBuilder::new_with_metadata(file, arrow_metadata); + + // Create a RowSelection that: + // - Selects rows 20-30 (in row group 0) + // - Selects rows 60-70 (in row group 1) + // - Selects rows 120-130 (in row group 2) + let selection = RowSelection::from(vec![ + // Skip first 20 rows + parquet::arrow::arrow_reader::RowSelector::skip(20), + // Select rows 20-30 + parquet::arrow::arrow_reader::RowSelector::select(10), + // Skip rows 30-60 + parquet::arrow::arrow_reader::RowSelector::skip(30), + // Select rows 60-70 + parquet::arrow::arrow_reader::RowSelector::select(10), + // Skip rows 70-120 + parquet::arrow::arrow_reader::RowSelector::skip(50), + // Select rows 120-130 + parquet::arrow::arrow_reader::RowSelector::select(10), + ]); + + let reader = builder + .with_row_selection(selection) + // make sure we're actually using the page index and not using a mask + .with_row_selection_policy(RowSelectionPolicy::Selectors) + .build() + .unwrap(); + + // Collect all batches + let batches: Vec<RecordBatch> = reader.collect::<Result<Vec<_>, _>>().unwrap(); + + // Step 7: Check statistics + // Note: We need to get the provider reference from the metadata to check stats + let provider_ref = metadata_with_custom_index + .page_index() + .unwrap() + .as_any() + .downcast_ref::<SelectivePageIndexProvider>() + .expect("Expected SelectivePageIndexProvider"); + + // Check that the custom provider was used + let (hits, misses) = provider_ref.stats(); + assert!(hits > 0, "provider should have hits"); + assert!(misses > 0, "provider should have misses"); + + // Since we didn't use predicates, the column index should be untouched + let (col_hits, col_misses) = provider_ref.column_index_stats(); + assert_eq!(col_hits, 0, "column index should not be used"); + assert_eq!(col_misses, 0, "column index should not be used"); + + let (off_hits, off_misses) = provider_ref.offset_index_stats(); + assert!(off_hits > 0, "offset index should have hits"); + assert!(off_misses > 0, "offset index should have misses"); + + // Step 8: Verify correct results + // We expect 30 rows total: 10 from each selected range + let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum(); + assert_eq!(total_rows, 30); + + // Verify the actual data values + let mut all_ids = Vec::new(); + let mut all_values = Vec::new(); + let mut all_names = Vec::new(); + let mut all_scores = Vec::new(); + + for batch in batches { + let ids = batch + .column(0) + .as_any() + .downcast_ref::<Int32Array>() + .unwrap(); + let values = batch + .column(1) + .as_any() + .downcast_ref::<Int32Array>() + .unwrap(); + let names = batch + .column(2) + .as_any() + .downcast_ref::<StringArray>() + .unwrap(); + let scores = batch + .column(3) + .as_any() + .downcast_ref::<Int32Array>() + .unwrap(); + + all_ids.extend(ids.iter().flatten()); + all_values.extend(values.iter().flatten()); + all_names.extend(names.iter().flatten().map(|s| s.to_owned())); + all_scores.extend(scores.iter().flatten()); + } + + // Expected values for the selected rows: + // rows 20-30: id=20..30, value=40..60, score=60..90 + // rows 60-70: id=60..70, value=120..140, score=180..210 + // rows 120-130: id=120..130, value=240..260, score=360..390 + let expected_ids: Vec<i32> = (20..30).chain(60..70).chain(120..130).collect(); + let expected_values: Vec<i32> = expected_ids.iter().map(|x| x * 2).collect(); + let expected_scores: Vec<i32> = expected_ids.iter().map(|x| x * 3).collect(); + let expected_names: Vec<String> = expected_ids.iter().map(|x| format!("name{x}")).collect(); + + assert_eq!(all_ids, expected_ids, "IDs don't match expected values"); + assert_eq!( + all_values, expected_values, + "Values don't match expected values" + ); + assert_eq!( + all_scores, expected_scores, + "Scores don't match expected values" + ); + assert_eq!( + all_names, expected_names, + "Names don't match expected values" + ); +} + +/// A custom PageIndexProvider that only stores indexes for a subset of columns +/// +/// This provider mimics a real-world scenario where an application has loaded +/// page indexes selectively based on query predicates and projections. +#[derive(Debug)] +struct SelectivePageIndexProvider { + file_bytes: Bytes, + // indexes are accessed first by row_group index and then by column index + column_indexes: Option<HashMap<usize, HashMap<usize, ColumnIndexMetaData>>>, + offset_indexes: Option<HashMap<usize, HashMap<usize, OffsetIndexMetaData>>>, + // Usage statistics + column_index_hits: Arc<AtomicUsize>, + column_index_misses: Arc<AtomicUsize>, + offset_index_hits: Arc<AtomicUsize>, + offset_index_misses: Arc<AtomicUsize>, +} + +impl SelectivePageIndexProvider { + fn new(file_bytes: Bytes) -> Self { + Self { + file_bytes, + column_indexes: None, + offset_indexes: None, + column_index_hits: Arc::new(AtomicUsize::new(0)), + column_index_misses: Arc::new(AtomicUsize::new(0)), + offset_index_hits: Arc::new(AtomicUsize::new(0)), + offset_index_misses: Arc::new(AtomicUsize::new(0)), + } + } + + /// Get statistics about column index usage + pub fn column_index_stats(&self) -> (usize, usize) { + ( + self.column_index_hits.load(Ordering::Relaxed), + self.column_index_misses.load(Ordering::Relaxed), + ) + } + + /// Get statistics about offset index usage + pub fn offset_index_stats(&self) -> (usize, usize) { + ( + self.offset_index_hits.load(Ordering::Relaxed), + self.offset_index_misses.load(Ordering::Relaxed), + ) + } + + /// Get combined statistics as (hits, misses) for both index types + pub fn stats(&self) -> (usize, usize) { + let (col_hits, col_misses) = self.column_index_stats(); + let (off_hits, off_misses) = self.offset_index_stats(); + (col_hits + off_hits, col_misses + off_misses) + } + + /// Reset all statistics counters to zero + pub fn reset_stats(&self) { + self.column_index_hits.store(0, Ordering::Relaxed); + self.column_index_misses.store(0, Ordering::Relaxed); + self.offset_index_hits.store(0, Ordering::Relaxed); + self.offset_index_misses.store(0, Ordering::Relaxed); + } + + /// Fetch and parse the column index for the given row group and column + fn fetch_column_index( + &mut self, + row_group_idx: usize, + column_idx: usize, + metadata: &ParquetMetaData, + ) -> parquet::errors::Result<()> { + let map = self.column_indexes.get_or_insert_with(HashMap::new); + let rg = map.entry(row_group_idx).or_default(); + if let Entry::Vacant(e) = rg.entry(column_idx) { + let column = metadata.row_group(row_group_idx).column(column_idx); + let range = column.column_index_range(); + if let Some(range) = range { + let idx_bytes = self + .file_bytes + .slice(range.start as usize..range.end as usize); + let idx = decode_column_index(&idx_bytes, column.column_type())?; + e.insert(idx); + } + } + Ok(()) + } + + /// Fetch and parse the offset index for the given row group and column + fn fetch_offset_index( + &mut self, + row_group_idx: usize, + column_idx: usize, + metadata: &ParquetMetaData, + ) -> parquet::errors::Result<()> { + let map = self.offset_indexes.get_or_insert_with(HashMap::new); + let rg = map.entry(row_group_idx).or_default(); + if let Entry::Vacant(e) = rg.entry(column_idx) { + let column = metadata.row_group(row_group_idx).column(column_idx); + let range = column.offset_index_range(); + if let Some(range) = range { + let idx_bytes = self + .file_bytes + .slice(range.start as usize..range.end as usize); + let idx = decode_offset_index(&idx_bytes)?; + e.insert(idx); + } + } + Ok(()) + } +} + +impl PageIndexProvider for SelectivePageIndexProvider { + fn has_offset_indexes(&self) -> bool { + self.offset_indexes.is_some() + } + + fn has_column_indexes(&self) -> bool { + self.column_indexes.is_some() + } + + fn column_index( + &self, + row_group_idx: usize, + column_idx: usize, + ) -> Option<&ColumnIndexMetaData> { + let result = self + .column_indexes + .as_ref()? + .get(&row_group_idx)? + .get(&column_idx); + + if result.is_some() { + self.column_index_hits.fetch_add(1, Ordering::Relaxed); + } else { + self.column_index_misses.fetch_add(1, Ordering::Relaxed); + } + + result + } + + fn offset_index( + &self, + row_group_idx: usize, + column_idx: usize, + ) -> Option<&OffsetIndexMetaData> { + let result = self + .offset_indexes + .as_ref()? + .get(&row_group_idx)? + .get(&column_idx); + + if result.is_some() { + self.offset_index_hits.fetch_add(1, Ordering::Relaxed); + } else { + self.offset_index_misses.fetch_add(1, Ordering::Relaxed); + } + + result + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } +} + +/// Create a test parquet file with multiple row groups and multiple pages per column +fn create_test_file() -> NamedTempFile { + let temp_file = tempfile::Builder::new() + .prefix("custom_page_index_test") + .suffix(".parquet") + .tempfile() + .expect("tempfile creation"); Review Comment: **Nit (B5).** - File I/O: this PR still writes a `NamedTempFile` and reads it back three times (L51, L54, L108). Commit `b682ebecff` ("rework tests to do no file i/o") on the [#11157](https://github.com/apache/arrow-rs/pull/11157) branch already changes this test to `Bytes`. Move that commit into this PR, so that this PR does not need a follow-up. - One 150-row batch gives the same layout (3 row groups, 5 pages per chunk, checked), because `set_max_row_group_row_count(Some(50))` splits it. The nested loop and `flush()` are not necessary. - Use one `fn make_batch(ids: &[i32]) -> RecordBatch` for the writer and for the check. Then L163-223 become: ```rust let ids: Vec<i32> = (20..30).chain(60..70).chain(120..130).collect(); let expected = make_batch(&ids); assert_eq!(concat_batches(&expected.schema(), &batches).unwrap(), expected); ``` (`concat_batches` is `arrow::compute::concat_batches`.) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
