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]

Reply via email to