adriangb opened a new pull request, #11238:
URL: https://github.com/apache/arrow-rs/pull/11238

   > [!NOTE]
   > This PR is stacked on https://github.com/apache/arrow-rs/pull/11218 (test) 
and https://github.com/apache/arrow-rs/pull/11235 (`release_ranges`). Only the 
last 2 commits are new: +2,643 lines, about 1,300 of them tests. I will rebase 
when those PRs merge.
   
   # Which issue does this PR close?
   
   - Part of https://github.com/apache/arrow-rs/issues/11234 (step E).
   
   # Rationale for this change
   
   Today `ParquetPushDecoder` decodes nothing until all bytes that a row group 
reads are pushed. https://github.com/apache/arrow-rs/pull/11218 shows this. For 
large row groups on object storage, the time to the first batch increases with 
the size of 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 and the first data page of each 
column
   let DecodeResult::NeedsData(ranges) = decoder.try_decode()? else { panic!() 
};
   decoder.push_ranges(ranges, fetch(&ranges).await?)?;
   
   // The first batch is ready. The rest of the row group was not requested.
   let DecodeResult::Data(batch) = decoder.try_decode()? else { panic!() };
   ```
   
   This PR is the smallest complete version. It works with and without a 
`RowFilter`, and has no temporary errors or fallbacks. Later PRs in 
https://github.com/apache/arrow-rs/issues/11234 add optimizations:
   
   | Not in this PR | Behavior in this PR | Later PR |
   |---|---|---|
   | Release each data page when all readers have passed it | The decoder holds 
the pages of a row group until the row group ends (the same peak as `RowGroup`) 
| F |
   | Predicate cache | The output decodes the predicate columns that it reads 
again | G |
   | Configured `RowSelectionPolicy` in each window | Each window uses 
selectors | H |
   
   # What changes are included in this PR?
   
   ## New public API
   
   | Item | Purpose |
   |---|---|
   | `FetchGranularity::{RowGroup, Batch}` (`#[non_exhaustive]`, default 
`RowGroup`) | How `try_decode` requests data and decodes |
   | `ParquetPushDecoderBuilder::with_fetch_granularity` | Sets it |
   
   The default does not change.
   
   ## How it works
   
   The design is in the module documentation of 
`reader_builder/incremental.rs`. In short:
   
   ```text
              push                 ingest (move)               get
    caller ─────────▶ PushBuffers ──────────────▶ PageStore ◀─────── column 
readers
   ```
   
   | Part | Where |
   |---|---|
   | `PageStore` and `ColumnChunkData::Shared`: column chunk data that the 
decoder can add pages to while a reader uses it | `push_decoder/page_store.rs`, 
`in_memory_row_group.rs` |
   | `IncrementalRowGroup`: one reader per predicate and one for the output, 
kept for the row group. A state machine (`Idle` → `Predicate` → `Output`) 
filters one window of `batch_size` rows at a time. Without predicates, a window 
goes to the output queue without I/O. | 
`push_decoder/reader_builder/incremental.rs` |
   | `ParquetRecordBatchReader::into_array_reader` (crate-private): drive one 
set of column readers with a sequence of short read plans | 
`arrow_reader/mod.rs` |
   | Wiring into the row-group state machine and the decoder | 
`reader_builder/mod.rs`, `remaining.rs`, `push_decoder/mod.rs` |
   
   ## Behavior with `FetchGranularity::Batch`
   
   | Topic | Behavior |
   |---|---|
   | `NeedsData` | Asks only for the pages that the next step reads, including 
dictionary pages. These are the same ranges that `RowGroup` requests, split per 
step. |
   | Decoding | Each reader decodes each page and each dictionary one time |
   | Batches | The same batches, in the same order, as `RowGroup` |
   | Memory | The pages of a row group are held until the row group ends. Then 
all bytes of the row group are released, as with `RowGroup` after 
https://github.com/apache/arrow-rs/pull/11235. |
   | No offset index | The column is requested as one column chunk |
   | Row-group boundary | Visible after the last batch of a row group, so 
`into_builder` works |
   | `try_next_reader` | Not changed. In a row group, only the method that 
started it can continue it. The other method returns an error and does not 
consume the decoder. |
   | `batch_size` of 0 | `build` returns an error if the file has rows |
   
   # Are these changes tested?
   
   Yes.
   
   - `test_decoder_first_pages_only_batch_granularity`: the counterpart of the 
test from https://github.com/apache/arrow-rs/pull/11218, on the same file.
   - `incremental_tests.rs` compares both granularities (same batches) on two 
files: one where all columns share page boundaries, and one with different page 
boundaries per column and uneven row groups. It covers batch sizes around page 
and row-group sizes, list and struct projections, row selections that skip 
pages, offset and limit, files without an offset index, predicate chains (also 
on nested columns and with null results), and randomized combinations of all of 
these. Other tests check the request sizes, that no range is requested two 
times, the row-group boundary, `into_builder`, and `try_next_reader` in both 
directions.
   - Unit tests for `PageStore` and the row-range helpers.
   
   # Are there any user-facing changes?
   
   Yes, a new opt-in API: `FetchGranularity` and 
`ParquetPushDecoderBuilder::with_fetch_granularity`. There are no breaking 
changes.
   
   # Open questions
   
   1. **API shape.** An enum (this PR), a row window size (for example 
`with_fetch_rows(n)`), or a separate method?
   2. **Predicate batches.** With a `RowFilter`, a predicate gets batches of at 
most `batch_size` rows from one window. The output is the same as with 
`RowGroup`, but a predicate that keeps state across calls can see different 
batch boundaries. Is this acceptable?
   3. **Row groups without output rows.** Such a row group passes without a 
visible row-group boundary. Must the decoder stop there 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]

Reply via email to