parquet: Add end-to-end test of custom PageIndexProvider - #11132
Conversation
adriangb
left a comment
There was a problem hiding this comment.
@etseidl AI disclaimer: this review text is largely AI generated, although I worked with my agent to guide the style of review, things to look for and how to present them as well as discussing some of the individual findings (but not all) and reading over the final review to make sure it makes sense at a high level. If this is not helpful please let me know and I will change strategies.
This test reads only with the sync reader. The same provider and RowSelection fail with the async reader and the push decoder, because of a bug that is already on main. This PR does not cause the bug and does not make it worse, so it is not a blocker. The fix is in #11182. When it is on main, this test can also run with the async reader and the push decoder.
Impact (PR head b7b529de8a, same file, provider and selection as this test, scratch test not committed):
| Reader | Result |
|---|---|
ParquetRecordBatchReaderBuilder (sync) |
30 rows, correct |
ParquetRecordBatchStreamBuilder (async) |
Err(General("Invalid column index 2, column was not fetched")) |
ParquetPushDecoderBuilder (push) |
Err(General("Invalid column index 2, column was not fetched")) |
Current main (e2db103480) gives the same result. With the fix in #11182, all three readers return the correct 30 rows.
| id | severity | summary | where |
|---|---|---|---|
| B1 | Bug (pre-existing), Test gap | The test uses only the sync reader. Async and push fail for the same scenario. | inline L130-138, body |
| B2 | Test gap | Chunks without an index always come after chunks with an index. An incomplete fix passes this layout. | inline L70-81 |
| B3 | Test gap | The stats asserts pass when the reader ignores the offset index. | inline L149-161 |
| B4 | Nit | The provider helper can be much smaller. | inline L231-241 |
| B5 | Nit | File I/O, writer loop and result checks can be smaller. | inline L388-394 |
| B6 | Nit | Full RowSelector paths; metadata load through a reader builder. |
body |
B1: pre-existing bug in InMemoryRowGroup (not caused by this PR)
The async reader and the push decoder use InMemoryRowGroup::fetch_ranges and fill_column_chunks. These functions align page_start_offsets to the columns by position. A chunk without an offset index gets a whole-chunk range, but no page_start_offsets entry.
flowchart LR
F["fetch_ranges, rg 0<br/>col 0, 1: page ranges + entry<br/>col 2, 3: whole chunk, no entry"] --> O["page_start_offsets = [col 0, col 1]"]
O --> L["fill_column_chunks<br/>col 0 = entry 0, col 1 = entry 1<br/>col 2, 3 = nothing"]
L --> E["column_chunks(2)<br/>Invalid column index 2, column was not fetched"]
-
Code:
parquet/src/arrow/in_memory_row_group.rsL88-92 (no entry forNone) and L161 (page_start_offsets.next()). -
The bug is in release 60.0.0 and later, since #10719 (
5ec9eaf033). Other triggers: a valid file where only some chunks have an offset index (no custom provider), and aRowFilterwith noRowSelection(the fetch after the predicate has a selection). I checked theRowFiltertrigger with this provider: async gives the same error. -
The sync reader gets
page_locationsfor each column separately (ReaderPageIterator::next_page_reader). It does not usefetch_ranges. This is why this test passes. -
Fix: #11182. A chunk without an offset index gets a
Noneentry and aDensechunk, so the positions stay aligned. Its tests cover a missing offset index before, between and after indexed columns, and row groups without any offset index. Each test uses the sync, async and push readers, with aRowSelectionand with aRowFilter.
Suggested order (not a blocker):
- This PR can merge before or after the fix. It does not regress anything.
- I prefer that the fix (#11182) merges first. Then this test can run with all three readers (snippet at L130-138), and it also protects the new mask paths in #11157.
- If this PR merges first, add the async and push variants with
#[ignore = "https://github.com/apache/arrow-rs/pull/11182"], or add them in the fix PR.
B6: other nits
- L115-128: import
RowSelectorand remove the sixparquet::arrow::arrow_reader::prefixes. - L53-61:
ParquetMetaDataReader::new().parse_and_finish(&file_bytes)gives the metadata without a reader builder. Its default policy isSkip.
Co-authored-by: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com>
|
Implemented all suggestions...it got a little bigger 😮 (TBF it is testing more now). |
…ex (#11182) # Which issue does this PR close? - None. There is no separate issue. This PR is related to #11132 and #11157. I found the bug while I reviewed them. # Rationale for this change The async reader (`ParquetRecordBatchStreamBuilder`) and the push decoder (`ParquetPushDecoderBuilder`) fail on a valid read. The sync reader (`ParquetRecordBatchReaderBuilder`) reads the same data correctly. Example: `file` is a valid Parquet file with the columns `a`, `b` and `c` (400 rows, 2 row groups, 50 rows per page). Columns `a` and `c` have an offset index. Column `b` does not have an offset index. ```rust let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Optional); let selection = RowSelection::from(vec![RowSelector::skip(150), RowSelector::select(100)]); // Sync reader: OK, 100 rows let batches = ParquetRecordBatchReaderBuilder::try_new_with_options(file.clone(), options.clone())? .with_row_selection(selection.clone()) .build()? .collect::<Result<Vec<_>, _>>()?; // Async reader: Err(General("Invalid column index 2, column was not fetched")) let batches: Vec<RecordBatch> = ParquetRecordBatchStreamBuilder::new_with_options(std::io::Cursor::new(file), options) .await? .with_row_selection(selection) .build()? .try_collect() .await?; ``` <details><summary>How to make <code>file</code></summary> ```rust use arrow_array::{ArrayRef, Int32Array, RecordBatch}; use bytes::Bytes; use parquet::arrow::ArrowWriter; use parquet::column::writer::ColumnCloseResult; use parquet::file::metadata::{PageIndexPolicy, ParquetMetaDataReader}; use parquet::file::properties::WriterProperties; use parquet::file::writer::SerializedFileWriter; use std::sync::Arc; fn make_file() -> Bytes { // A file with a full page index let column = |start| Arc::new(Int32Array::from_iter_values(start..start + 400)) as ArrayRef; let batch = RecordBatch::try_from_iter([("a", column(0)), ("b", column(400)), ("c", column(800))]) .unwrap(); let props = WriterProperties::builder() .set_max_row_group_row_count(Some(200)) .set_data_page_row_count_limit(50) .set_write_batch_size(50) .build(); let mut buf = Vec::new(); let mut writer = ArrowWriter::try_new(&mut buf, batch.schema(), Some(props)).unwrap(); writer.write(&batch).unwrap(); writer.close().unwrap(); let source = Bytes::from(buf); // Copy the column chunks to a new file, without the page index of column `b` let metadata = ParquetMetaDataReader::new() .with_page_index_policy(PageIndexPolicy::Required) .parse_and_finish(&source) .unwrap(); let page_index = metadata.page_index().unwrap(); let schema = metadata.file_metadata().schema_descr().root_schema_ptr(); let mut buf = Vec::new(); let mut writer = SerializedFileWriter::new(&mut buf, schema, Default::default()).unwrap(); for (rg, rg_meta) in metadata.row_groups().iter().enumerate() { let mut rg_writer = writer.next_row_group().unwrap(); for (col, col_meta) in rg_meta.columns().iter().enumerate() { let keep = col != 1; let close = ColumnCloseResult { bytes_written: col_meta.compressed_size() as u64, rows_written: rg_meta.num_rows() as u64, metadata: col_meta.clone(), bloom_filter: None, column_index: page_index.column_index(rg, col).filter(|_| keep).cloned(), offset_index: page_index.offset_index(rg, col).filter(|_| keep).cloned(), }; rg_writer.append_column(&source, close).unwrap(); } rg_writer.close().unwrap(); } writer.close().unwrap(); Bytes::from(buf) } ``` </details> | Reader | Before (`main`, 60.0.0) | After (this PR) | |---|---|---| | Sync (`ParquetRecordBatchReaderBuilder`) | 100 rows | 100 rows | | Async (`ParquetRecordBatchStreamBuilder`) | Error: `Invalid column index 2, column was not fetched` | 100 rows | | Push decoder (`ParquetPushDecoderBuilder`) | Error: `Invalid column index 2, column was not fetched` | 100 rows | The error occurs when all of these conditions are true: 1. The reader loads the page index (`PageIndexPolicy::Optional` or `Required`). 2. In a row group, some of the column chunks to fetch have an offset index, and some do not. 3. The read has a `RowSelection` or a `RowFilter`. After a predicate, the reader always fetches the remaining columns with a selection. Condition 2 occurs with a valid file in which only some column chunks have an offset index, or with a custom `PageIndexProvider`. The offset index mask in #11157 makes condition 2 a usual case. Release 60.0.0 is affected. The bug started with #10719 (commit 5ec9eaf033), which made the offset index an `Option` for each column chunk. Release 59.3.0 is not affected. ## Cause `InMemoryRowGroup::fetch_ranges` fetches the full column chunk when the chunk has no offset index, but it does not push an entry to `page_start_offsets`. `fill_column_chunks` reads `page_start_offsets` by position. Thus the entries go to the wrong columns, and the last column gets no data. | Column | Offset index | `page_start_offsets` entry | Before this PR | After this PR | |---|---|---|---|---| | `a` | Yes | `offsets_a` | `a` gets `offsets_a` | `Sparse` with `offsets_a` | | `b` | No | Missing (now `None`) | `b` gets `offsets_c` (wrong) | `Dense` (full chunk) | | `c` | Yes | `offsets_c` | `c` gets no data: error | `Sparse` with `offsets_c` | # What changes are included in this PR? - `page_start_offsets` is now `Option<Vec<Option<Vec<u64>>>>` in `FetchRanges`, `fill_column_chunks` and the push decoder `DataRequest`. - `fetch_ranges` pushes `None` for a column chunk without an offset index. - `fill_column_chunks` stores a `None` entry as `ColumnChunkData::Dense` (the full chunk). It stores a `Some(offsets)` entry as `ColumnChunkData::Sparse`, as before. A `Sparse` chunk with one range at the chunk start does not work. Without page locations, `SerializedPageReader` reads each page at its own offset, and `Sparse` accepts only an exact page start. That change gives the error `Invalid offset in sparse column chunk data: ..., no matching page found`. # Are these changes tested? Yes. The new file `parquet/tests/arrow_reader/partial_offset_index.rs` makes a valid file in memory in which only some column chunks have an offset index. It uses `SerializedRowGroupWriter::append_column` for this. Each test reads rows 150..250 with the sync reader (control), the async reader and the push decoder. Each test uses a `RowSelection`, then a `RowFilter`, and compares the output with the expected rows. | Test | Column chunks without an offset index | Async reader and push decoder on `main` | |---|---|---| | `test_no_offset_index_first_column` | `a` (before the indexed columns) | Fail with `RowSelection` | | `test_no_offset_index_middle_column` | `b` (between the indexed columns) | Fail with `RowSelection` and with `RowFilter` | | `test_no_offset_index_last_column` | `c` (after the indexed columns) | Fail with `RowSelection` and with `RowFilter` | | `test_no_offset_index_row_group` | All columns in row group 1 | Fail with `RowSelection` and with `RowFilter` | All 4 tests pass with this PR. With only a `RowFilter` on `a`, the first-column case passes on `main`, because the predicate step already fetched column `a`. # Are there any user-facing changes? No API changes. Reads that failed with `Invalid column index N, column was not fetched` now return the correct rows. Note on AI use: Claude Code wrote the fix and the tests. The author reviewed them. 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
I'll take token approval anyday ! Especially when that means the PR is already in good shape |
| let mut metadata_builder = metadata.clone().into_builder(); | ||
| metadata_builder = metadata_builder.set_page_index(Some(Arc::new(provider))); |
There was a problem hiding this comment.
a minor usability thing -- it would be nice to have a with type method, so this could be
let metadata_builder = metadata.clone()
.into_builder()
.with_page_index(Some(Arc::new(provider)));But I don't think it is necessary
There was a problem hiding this comment.
Like this? 🤣 Forgot that set_page_index could be chained.
let metadata_builder = metadata
.clone()
.into_builder()
.set_page_index(Some(provider.clone()));| // - Selects rows 120-130 (in row group 2) | ||
| let selection = RowSelection::from(vec![ | ||
| // Skip first 20 rows | ||
| parquet::arrow::arrow_reader::RowSelector::skip(20), |
There was a problem hiding this comment.
can we add a use parquet::arrow::arrow_reader::RowSelector to avoid the repetition here?
| } | ||
|
|
||
| // Note: We need to get the provider reference from the metadata to check stats | ||
| let provider_ref = metadata_with_custom_index |
There was a problem hiding this comment.
It might be clearer to just keep this reference when it was first created, rather than trying to get it again
| let provider = SelectivePageIndexProvider::new( | ||
| file_bytes.clone(), | ||
| metadata, | ||
| &[(0, 0), (0, 1), (2, 0), (2, 2)], |
There was a problem hiding this comment.
this is documented later on, but it would help me read this example if we could explain what these numbers represented
something like
// values are (row_group index, column index)
&[(0, 0), (0, 1), (2, 0), (2, 2)],) # Which issue does this PR close? - Closes apache#11029. # Rationale for this change Add test to demonstrate that read functionality is preserved when using a custom `PageIndexProvider`. # What changes are included in this PR? Adds test # Are these changes tested? N/A # Are there any user-facing changes? No # AI Usage The test was written by Claude Code and lightly edited by me --------- Co-authored-by: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com>
…ex (apache#11182) # Which issue does this PR close? - None. There is no separate issue. This PR is related to apache#11132 and apache#11157. I found the bug while I reviewed them. # Rationale for this change The async reader (`ParquetRecordBatchStreamBuilder`) and the push decoder (`ParquetPushDecoderBuilder`) fail on a valid read. The sync reader (`ParquetRecordBatchReaderBuilder`) reads the same data correctly. Example: `file` is a valid Parquet file with the columns `a`, `b` and `c` (400 rows, 2 row groups, 50 rows per page). Columns `a` and `c` have an offset index. Column `b` does not have an offset index. ```rust let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Optional); let selection = RowSelection::from(vec![RowSelector::skip(150), RowSelector::select(100)]); // Sync reader: OK, 100 rows let batches = ParquetRecordBatchReaderBuilder::try_new_with_options(file.clone(), options.clone())? .with_row_selection(selection.clone()) .build()? .collect::<Result<Vec<_>, _>>()?; // Async reader: Err(General("Invalid column index 2, column was not fetched")) let batches: Vec<RecordBatch> = ParquetRecordBatchStreamBuilder::new_with_options(std::io::Cursor::new(file), options) .await? .with_row_selection(selection) .build()? .try_collect() .await?; ``` <details><summary>How to make <code>file</code></summary> ```rust use arrow_array::{ArrayRef, Int32Array, RecordBatch}; use bytes::Bytes; use parquet::arrow::ArrowWriter; use parquet::column::writer::ColumnCloseResult; use parquet::file::metadata::{PageIndexPolicy, ParquetMetaDataReader}; use parquet::file::properties::WriterProperties; use parquet::file::writer::SerializedFileWriter; use std::sync::Arc; fn make_file() -> Bytes { // A file with a full page index let column = |start| Arc::new(Int32Array::from_iter_values(start..start + 400)) as ArrayRef; let batch = RecordBatch::try_from_iter([("a", column(0)), ("b", column(400)), ("c", column(800))]) .unwrap(); let props = WriterProperties::builder() .set_max_row_group_row_count(Some(200)) .set_data_page_row_count_limit(50) .set_write_batch_size(50) .build(); let mut buf = Vec::new(); let mut writer = ArrowWriter::try_new(&mut buf, batch.schema(), Some(props)).unwrap(); writer.write(&batch).unwrap(); writer.close().unwrap(); let source = Bytes::from(buf); // Copy the column chunks to a new file, without the page index of column `b` let metadata = ParquetMetaDataReader::new() .with_page_index_policy(PageIndexPolicy::Required) .parse_and_finish(&source) .unwrap(); let page_index = metadata.page_index().unwrap(); let schema = metadata.file_metadata().schema_descr().root_schema_ptr(); let mut buf = Vec::new(); let mut writer = SerializedFileWriter::new(&mut buf, schema, Default::default()).unwrap(); for (rg, rg_meta) in metadata.row_groups().iter().enumerate() { let mut rg_writer = writer.next_row_group().unwrap(); for (col, col_meta) in rg_meta.columns().iter().enumerate() { let keep = col != 1; let close = ColumnCloseResult { bytes_written: col_meta.compressed_size() as u64, rows_written: rg_meta.num_rows() as u64, metadata: col_meta.clone(), bloom_filter: None, column_index: page_index.column_index(rg, col).filter(|_| keep).cloned(), offset_index: page_index.offset_index(rg, col).filter(|_| keep).cloned(), }; rg_writer.append_column(&source, close).unwrap(); } rg_writer.close().unwrap(); } writer.close().unwrap(); Bytes::from(buf) } ``` </details> | Reader | Before (`main`, 60.0.0) | After (this PR) | |---|---|---| | Sync (`ParquetRecordBatchReaderBuilder`) | 100 rows | 100 rows | | Async (`ParquetRecordBatchStreamBuilder`) | Error: `Invalid column index 2, column was not fetched` | 100 rows | | Push decoder (`ParquetPushDecoderBuilder`) | Error: `Invalid column index 2, column was not fetched` | 100 rows | The error occurs when all of these conditions are true: 1. The reader loads the page index (`PageIndexPolicy::Optional` or `Required`). 2. In a row group, some of the column chunks to fetch have an offset index, and some do not. 3. The read has a `RowSelection` or a `RowFilter`. After a predicate, the reader always fetches the remaining columns with a selection. Condition 2 occurs with a valid file in which only some column chunks have an offset index, or with a custom `PageIndexProvider`. The offset index mask in apache#11157 makes condition 2 a usual case. Release 60.0.0 is affected. The bug started with apache#10719 (commit apache@5ec9eaf033), which made the offset index an `Option` for each column chunk. Release 59.3.0 is not affected. ## Cause `InMemoryRowGroup::fetch_ranges` fetches the full column chunk when the chunk has no offset index, but it does not push an entry to `page_start_offsets`. `fill_column_chunks` reads `page_start_offsets` by position. Thus the entries go to the wrong columns, and the last column gets no data. | Column | Offset index | `page_start_offsets` entry | Before this PR | After this PR | |---|---|---|---|---| | `a` | Yes | `offsets_a` | `a` gets `offsets_a` | `Sparse` with `offsets_a` | | `b` | No | Missing (now `None`) | `b` gets `offsets_c` (wrong) | `Dense` (full chunk) | | `c` | Yes | `offsets_c` | `c` gets no data: error | `Sparse` with `offsets_c` | # What changes are included in this PR? - `page_start_offsets` is now `Option<Vec<Option<Vec<u64>>>>` in `FetchRanges`, `fill_column_chunks` and the push decoder `DataRequest`. - `fetch_ranges` pushes `None` for a column chunk without an offset index. - `fill_column_chunks` stores a `None` entry as `ColumnChunkData::Dense` (the full chunk). It stores a `Some(offsets)` entry as `ColumnChunkData::Sparse`, as before. A `Sparse` chunk with one range at the chunk start does not work. Without page locations, `SerializedPageReader` reads each page at its own offset, and `Sparse` accepts only an exact page start. That change gives the error `Invalid offset in sparse column chunk data: ..., no matching page found`. # Are these changes tested? Yes. The new file `parquet/tests/arrow_reader/partial_offset_index.rs` makes a valid file in memory in which only some column chunks have an offset index. It uses `SerializedRowGroupWriter::append_column` for this. Each test reads rows 150..250 with the sync reader (control), the async reader and the push decoder. Each test uses a `RowSelection`, then a `RowFilter`, and compares the output with the expected rows. | Test | Column chunks without an offset index | Async reader and push decoder on `main` | |---|---|---| | `test_no_offset_index_first_column` | `a` (before the indexed columns) | Fail with `RowSelection` | | `test_no_offset_index_middle_column` | `b` (between the indexed columns) | Fail with `RowSelection` and with `RowFilter` | | `test_no_offset_index_last_column` | `c` (after the indexed columns) | Fail with `RowSelection` and with `RowFilter` | | `test_no_offset_index_row_group` | All columns in row group 1 | Fail with `RowSelection` and with `RowFilter` | All 4 tests pass with this PR. With only a `RowFilter` on `a`, the first-column case passes on `main`, because the predicate step already fetched column `a`. # Are there any user-facing changes? No API changes. Reads that failed with `Invalid column index N, column was not fetched` now return the correct rows. Note on AI use: Claude Code wrote the fix and the tests. The author reviewed them. 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
Which issue does this PR close?
Rationale for this change
Add test to demonstrate that read functionality is preserved when using a custom
PageIndexProvider.What changes are included in this PR?
Adds test
Are these changes tested?
N/A
Are there any user-facing changes?
No
AI Usage
The test was written by Claude Code and lightly edited by me