parquet: Add new PageIndex struct to encapsulate column and offset indexes - #10719
Conversation
49cbf86 to
7f47fac
Compare
|
Shal I give this one a look? |
Yes please. I think there's a todo left and it needs some better docs, but if you could take a look and see if you think this is the direction you wanted to head I'd appreciate it. |
|
Thank you -- I will do so but probably not until tomorrow (I need a clear mind) |
|
Some findings from a test upgrade Specifically this commit: apache/datafusion@419027d This pattern comes up a bunch (to see if the index is "complete" it has to test both column_indexes and offset_indexes. metadata
.page_index()
.is_some_and(|index| index.has_column_indexes() && index.has_offset_indexes())Maybe we could add a helper like PageIndex::is_complete that does the check and we could simplify to metadata
.page_index()
.is_some_and(PageIndex::is_complete)But otherwise the changes needed are pretty strightforward |
That's a good idea...there are a bunch of tests too that do the same thing. I'll get to that |
|
😍 |
|
Thanks for the review @alamb. I think this is ready to merge now. I'm going to get started on the builder to prototype partial population of the indexes. |
Replace uses of ParquetMetaData::column_index/offset_index and ParquetColumnIndex/ParquetOffsetIndex with the new PageIndex struct and its accessors, including PageIndex::is_complete.
|
werd! (as someone tell me the kids used to say) |
…indexes (apache#10719) # Which issue does this PR close? - Closes apache#8818. - Closes apache#10653 Note: this is stacked on apache#10653 # Rationale for this change Following up on apache#10653 (review) > If we are going to mess with the APIs I think it may be worth considering some more drastic changes # What changes are included in this PR? Try to hide some of the complexity of the page indexes behind a struct with accessors for access by row group, or row group and column. # Are these changes tested? Should be covered by existing # Are there any user-facing changes? Yes, this changes the public interface to the page indexes quite a bit.
…indexes (apache#10719) # Which issue does this PR close? - Closes apache#8818. - Closes apache#10653 Note: this is stacked on apache#10653 # Rationale for this change Following up on apache#10653 (review) > If we are going to mess with the APIs I think it may be worth considering some more drastic changes # What changes are included in this PR? Try to hide some of the complexity of the page indexes behind a struct with accessors for access by row group, or row group and column. # Are these changes tested? Should be covered by existing # Are there any user-facing changes? Yes, this changes the public interface to the page indexes quite a bit.
…indexes (apache#10719) # Which issue does this PR close? - Closes apache#8818. - Closes apache#10653 Note: this is stacked on apache#10653 # Rationale for this change Following up on apache#10653 (review) > If we are going to mess with the APIs I think it may be worth considering some more drastic changes # What changes are included in this PR? Try to hide some of the complexity of the page indexes behind a struct with accessors for access by row group, or row group and column. # Are these changes tested? Should be covered by existing # Are there any user-facing changes? Yes, this changes the public interface to the page indexes quite a bit.
…indexes (apache#10719) # Which issue does this PR close? - Closes apache#8818. - Closes apache#10653 Note: this is stacked on apache#10653 # Rationale for this change Following up on apache#10653 (review) > If we are going to mess with the APIs I think it may be worth considering some more drastic changes # What changes are included in this PR? Try to hide some of the complexity of the page indexes behind a struct with accessors for access by row group, or row group and column. # Are these changes tested? Should be covered by existing # Are there any user-facing changes? Yes, this changes the public interface to the page indexes quite a bit.
…indexes (apache#10719) # Which issue does this PR close? - Closes apache#8818. - Closes apache#10653 Note: this is stacked on apache#10653 # Rationale for this change Following up on apache#10653 (review) > If we are going to mess with the APIs I think it may be worth considering some more drastic changes # What changes are included in this PR? Try to hide some of the complexity of the page indexes behind a struct with accessors for access by row group, or row group and column. # Are these changes tested? Should be covered by existing # Are there any user-facing changes? Yes, this changes the public interface to the page indexes quite a bit.
…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>
…indexes (apache#10719) # Which issue does this PR close? - Closes apache#8818. - Closes apache#10653 Note: this is stacked on apache#10653 # Rationale for this change Following up on apache#10653 (review) > If we are going to mess with the APIs I think it may be worth considering some more drastic changes # What changes are included in this PR? Try to hide some of the complexity of the page indexes behind a struct with accessors for access by row group, or row group and column. # Are these changes tested? Should be covered by existing # Are there any user-facing changes? Yes, this changes the public interface to the page indexes quite a bit.
Which issue does this PR close?
ParquetMetaData#8818.Vec<Vec<Option<T>>>#10653Note: this is stacked on #10653
Rationale for this change
Following up on #10653 (review)
What changes are included in this PR?
Try to hide some of the complexity of the page indexes behind a struct with accessors for access by row group, or row group and column.
Are these changes tested?
Should be covered by existing
Are there any user-facing changes?
Yes, this changes the public interface to the page indexes quite a bit.