adriangb opened a new pull request, #11182: URL: https://github.com/apache/arrow-rs/pull/11182
# Which issue does this PR close? - None. There is no separate issue. This PR is related to https://github.com/apache/arrow-rs/pull/11132 and https://github.com/apache/arrow-rs/pull/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 https://github.com/apache/arrow-rs/pull/11157 makes condition 2 a usual case. Release 60.0.0 is affected. The bug started with https://github.com/apache/arrow-rs/pull/10719 (commit https://github.com/apache/arrow-rs/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) -- 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]
