adriangb commented on PR #11223: URL: https://github.com/apache/arrow-rs/pull/11223#issuecomment-5848201780
Automated review (Claude Fable 5.1). Findings on the 4 new commits only (ae76e574c4..dc8cddecaa). Method: read the new state machine end to end, then ran a second randomized equivalence suite (400 seeds) against `FetchGranularity::RowGroup` on a file the existing tests do not cover: columns with different page boundaries (28 / 18 / 140 / 35 pages per column in one row group), uneven row groups (700 / 700 / 400 rows), predicates on the list and struct columns, predicates that return nulls, up to 3 predicates, `with_row_groups` combined with a row selection, and both selection policies with and without the predicate cache. All of that matched, with no range requested twice. So the findings below are about edge input, performance and documentation, not about wrong rows. | # | Severity | Status | Where | Scenario | Suggested fix | |---|---|---|---|---|---| | 1 | bug | confirmed | [incremental.rs#L633](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/reader_builder/incremental.rs#L633) | `with_batch_size(0)` is accepted by the builder ([arrow_reader/mod.rs#L326](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/arrow_reader/mod.rs#L326) only clamps to `num_rows`). With `FetchGranularity::Batch` and a `RowFilter`, `try_decode` panics: `attempt to calculate the remainder with a divisor of zero`. Without a filter it returns `Internal Error: row group 0 ended after 0 of 700 rows` ([incremental.rs#L572](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/reader_builder/incremental.rs#L572)). `FetchGranularity::RowGroup` returns `Finished` for the same input. | Reject `batch_size == 0` in `ParquetPushDecoderBuilder::build` ([push_decoder/mod.rs#L401](https://github.com/a pache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/mod.rs#L401)) with a `ParquetError`, or clamp it to 1. Add the case to `incremental_tests.rs`. | | 2 | perf | confirmed | [incremental.rs#L527](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/reader_builder/incremental.rs#L527), [incremental.rs#L596](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/reader_builder/incremental.rs#L596), [incremental.rs#L1018](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/reader_builder/incremental.rs#L1018), [push_buffers.rs#L195](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/util/push_buffers.rs#L195) | Decode time is quadratic in the number of resident pushed buffers. `PushBuffers::release_range` rebuilds both `Vec`s on every call, and `has_range` is a linear scan. `ingest` calls `release_range` once per page it moves into the store, and the release paths call it once per page again. The docs recommend exac tly the usage that triggers this: push each page in its own buffer and push ahead. Measured (debug build, one row group, 4 Int64 columns, 50 rows per page, batch size 1000, whole row group pushed ahead one buffer per page): 4000 buffers 0.27 s, 8000 buffers 1.03 s, 16000 buffers 4.0 s. The same scans take 12 / 18 / 36 ms with `FetchGranularity::RowGroup`, and 22 / 47 / 104 ms with `Batch` when the caller pushes only what is requested. | Keep `PushBuffers` ordered by `range.start` (or index it with a `BTreeMap`) so `has_range` and `release_range` are `O(log n)`, or move every buffer that overlaps the active row group's column chunks into the `PageStore` in one pass when the row group starts and drop the per-page `release_range` calls. Add a test that pushes a row group ahead as per-page buffers and bounds the cost. | | 3 | docs / perf | confirmed | [push_decoder/mod.rs#L231](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/mod.rs#L231), [push_decoder/mod.rs#L256](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/mod.rs#L256) | With predicates, each window costs one `NeedsData` round trip per predicate plus one for the output, in sequence, because a later predicate's pages depend on the earlier one's result. For 1800 rows at batch size 100 the round trips are 18 / 28 / 44 for 0 / 1 / 2 predicates, against 3 / 6 / 9 with `RowGroup`. The docs say `NeedsData` "asks only for the bytes that the next batch reads" and do not say that a caller that pushes only what is requested pays `predicates + 1` sequential fetches per window on a latency-bound store. | State it in the `FetchGranularity::Batch` docs, and say that pushing ahead is the intended usage (which makes item 2 matter). | | 4 | docs | confirmed | [push_decoder/mod.rs#L256](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/mod.rs#L256), [incremental.rs#L617](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/reader_builder/incremental.rs#L617) | "decodes an output batch as soon as enough rows pass" is not exact. The decoder emits only when the queue holds more than `batch_size` rows, or when filtering is done. When exactly `batch_size` rows pass, it fetches and filters the next window first. This is by design (it must know the last batch of the row group), but the doc should say so. Also [push_decoder/mod.rs#L238](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/mod.rs#L238): at the end of a row group the decoder releases the chunks of the columns it reads, not "all remaining bytes of the row group's column chunks". Bytes pushed after `build` for columns outside the projection and the predicates stay resident until the decoder is dropped. | Adjust the two sentences. | | 5 | test gap | confirmed | [incremental_tests.rs#L105](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/incremental_tests.rs#L105) | The test file writes every column with 25 rows per page, so every column has the same page boundaries, and no random predicate touches a nested column or returns nulls. The page release logic (`release_passed_pages`, `SelectedRows`, `window_plan` with `Mask`) is per column and per row position, so different boundaries per column are the interesting case. I ran that case and it passes, so this is only coverage. | Add a second writer configuration to the randomized test: `set_data_page_size_limit(200)`, `set_write_batch_size(5)`, `set_data_page_row_count_limit(40)` gives 28 / 18 / 140 / 35 / 28 / 28 pages per column for the same schema. Add a predicate on `l` (`is_null`) and on `s.x`, and one that returns nulls. | | 6 | simplification | suspected | [incremental.rs#L610](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/reader_builder/incremental.rs#L610), [incremental.rs#L730](https://github.com/apache/arrow-rs/blob/dc8cddecaa011c8e9b3ac7f9106f4c69c5705555/parquet/src/arrow/push_decoder/reader_builder/incremental.rs#L730) | `let Mode::Filtered(state) = &mut self.mode else { unreachable!() }` appears 15 times, up to 5 times in one function, to satisfy the borrow checker. `fetch_ranges` builds a throwaway `InMemoryRowGroup` with `vec![None; num_columns]` on every step only to call `InMemoryRowGroup::fetch_ranges`. Not a bug, but the file is 1200 lines and this is where a reviewer will spend time. | Move the filtered state machine into `impl Filtered` with `&IncrementalConfig` and `&PageStore` parameters, and factor the range planning out of `InMemoryRowGroup::fetch_ranges` into a free function that takes the metadata and page index . | Checked and found no problem: batch boundaries versus `RowGroup` (all cases above), rows skipped by a selection when a page is only partly selected (whole-page skips read no bytes, partial skips land in a fetched page), dictionary pages and column chunks without an offset index stay resident until the row group ends, `RowFilter` stages with the predicate cache (windows are aligned to `batch_size`, so the cache batches of `CachedArrayReader` and the expanded fetch ranges match), nested columns (with an offset index the reader never reads past the requested record), offset and limit applied per window, zero-row row groups, `into_builder` and `release_unplanned_bytes`, `try_next_reader` while incremental, files with an offset index on only some columns (they fall back to whole column chunks per column), and no `unsafe` or reachable `unwrap` in the new code apart from item 1. 🤖 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]
