adriangb opened a new pull request, #11223: URL: https://github.com/apache/arrow-rs/pull/11223
> [!NOTE] > This PR is stacked on https://github.com/apache/arrow-rs/pull/11218. Only the 4 commits after alamb's commits are new. I will rebase when https://github.com/apache/arrow-rs/pull/11218 merges. # Which issue does this PR close? - Part of https://github.com/apache/arrow-rs/issues/6946. This follows the order agreed in https://github.com/apache/arrow-rs/issues/6946#issuecomment-5836698618: first, the push decoder emits a batch as soon as its data is pushed (this PR, starting from the test in https://github.com/apache/arrow-rs/pull/11218). The `ScanPlan` API (https://github.com/apache/arrow-rs/pull/10555) comes later and is not in this PR. # Rationale for this change Today `ParquetPushDecoder` decodes nothing until all projected bytes of a row group are pushed. https://github.com/apache/arrow-rs/pull/11218 shows this: with the first page of each column pushed, `try_decode` still returns `NeedsData` for the whole column chunks. For large row groups on object storage, the time to the first batch and the peak memory both grow with the row group. With this PR, a caller can opt in to batch granularity. The decoder then asks only for the pages of the next batch and returns the batch as soon as they are pushed: ```rust let mut decoder = ParquetPushDecoderBuilder::try_new_decoder(metadata)? .with_batch_size(100) .with_fetch_granularity(FetchGranularity::Batch) // new .build()?; // NeedsData asks for the dictionary page + first data page of each column let DecodeResult::NeedsData(ranges) = decoder.try_decode()? else { panic!() }; let data = fetch(&ranges).await?; decoder.push_ranges(ranges, data)?; // The first batch is ready. The rest of the row group was never requested. let DecodeResult::Data(batch) = decoder.try_decode()? else { panic!() }; ``` `test_decoder_first_pages_only_batch_granularity` (next to `test_decoder_first_pages_only` from https://github.com/apache/arrow-rs/pull/11218) tests exactly this. Numbers from a local simulation with the prototype of this change (one row group of 84 MB, 50 ms request latency). They are not a benchmark in this PR: | Mode | First batch | Peak buffered bytes | | --- | --- | --- | | `FetchGranularity::RowGroup` (default) | 451 ms | 84 MB | | `FetchGranularity::Batch`, caller read-ahead window of 8 MB | 89 ms | 8.1 MB | https://github.com/apache/datafusion/pull/24086 uses this mode. With it, DataFusion no longer needs its fallbacks for this case. # What changes are included in this PR? ## New public API | Item | Purpose | | --- | --- | | `FetchGranularity::{RowGroup, Batch}` (`#[non_exhaustive]`, default `RowGroup`) | Selects how `try_decode` requests data and decodes | | `ParquetPushDecoderBuilder::with_fetch_granularity` | Sets the mode | The default does not change. There is no other new public API. ## Behavior with `FetchGranularity::Batch` | Topic | Behavior | | --- | --- | | `NeedsData` | Asks only for the pages that the next batch reads, plus the dictionary pages they need | | Decoding | One set of column readers for the whole row group. No page is decoded twice. Dictionaries are decoded once. | | Memory | A data page is released after all readers pass its rows. The rest of the row group's bytes are released when the row group is done, including bytes that were pushed but not requested. | | `into_builder` / `build` | `build` releases pushed bytes that no row group of the decoder reads (for example, row groups skipped by a rebuild) | | Row-group boundary | Visible after the last batch of a row group, so `into_builder` works as with `try_next_reader` | | `try_next_reader` | Not changed. It still returns a reader for a whole row group. | | No offset index | The column is requested as one column chunk (row-group granularity for that column) | | `RowFilter` | Predicates run on one window of `batch_size` rows at a time. Output batches are the same as with `RowGroup`. | | Row selection policy | Each window uses the configured `RowSelectionPolicy` | ## Commits | Commit | Content | | --- | --- | | `feat(parquet): batch-granular fetching in ParquetPushDecoder` | The mode. Internal parts: `PageStore` + `ColumnChunkData::Shared` (pages added and removed while a reader is alive), crate-private `ParquetRecordBatchReader::into_array_reader`, `PushBuffers::release_range`, internal `page_spans` module, new `reader_builder/incremental.rs` state | | `test(parquet): cover FetchGranularity::Batch` | Counterpart of the test from https://github.com/apache/arrow-rs/pull/11218, and `incremental_tests.rs` | | `feat(parquet): release bytes of skipped row groups when a Batch decoder is built` | Release on `build`, with tests | | `perf(parquet): use the configured selection policy in batch windows` | Resolve the configured policy per window instead of always `Selectors` | # Are these changes tested? Yes. - `test_decoder_first_pages_only` from https://github.com/apache/arrow-rs/pull/11218 is not changed. It still asserts the default `RowGroup` behavior. - `test_decoder_first_pages_only_batch_granularity` shows the new behavior on the same file. - `incremental_tests.rs` compares both modes (same batches) for batch sizes around page and row-group sizes, list and struct projections, page-skipping row selections with each selection policy, offset and limit, files without an offset index, predicate chains with and without the predicate cache, and a randomized equivalence test. Other tests check request sizes, that dictionary pages and ranges are requested once, that bytes are released, the row-group boundary, and `into_builder` / `try_next_reader` interplay. # Are there any user-facing changes? Yes, new opt-in API: `FetchGranularity` and `ParquetPushDecoderBuilder::with_fetch_granularity`. There are no breaking changes. The default behavior does not change. # Open questions 1. **Opt-in or default?** This PR makes `Batch` opt-in. A middle option: by default, return a batch as soon as its data is present, but keep `NeedsData` asking for the whole row group. Callers that push whole row groups see no change, and callers that push pages early get batches early. 2. **API shape.** An enum (this PR), a row window size (for example `with_fetch_rows(n)`), or a separate method on the decoder? 3. **Predicate semantics.** With a `RowFilter`, the decoder evaluates predicates one window of `batch_size` rows at a time. The output is the same, but a predicate can get smaller batches than with `RowGroup`, and batch boundaries differ. A predicate that keeps state across calls can see this. Is this acceptable, or must the docs say more? 4. **Releasing caller-pushed bytes.** In `Batch` mode, the decoder releases bytes that the caller pushed but that no reader needs. The `into_builder` docs say that such bytes stay buffered for `RowGroup`. This PR documents the difference. Is it acceptable that the two modes differ here? 5. **Zero-row row groups.** A row group that produces no rows passes without a visible row-group boundary. Must the decoder stop at that boundary too? 🤖 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]
