adriangb opened a new pull request, #25727:
URL: https://github.com/apache/datafusion/pull/25727

   feat(parquet): adaptive placement of scan filter conjuncts
   
   ## Which issue does this PR close?
   
   - Part of https://github.com/apache/datafusion/issues/22883. Design notes: 
https://claude.ai/artifact/SSz7t6hPyhFWp1MDPecVqt
   - **Depends on https://github.com/apache/datafusion/pull/25682, 
https://github.com/apache/datafusion/pull/25722 and 
https://github.com/apache/datafusion/pull/25726.** Review only the commits in 
the "This PR" part of the commit table.
   - Prior art: https://github.com/apache/datafusion/pull/22236 
(`SelectivityTracker` cost model) and 
https://github.com/apache/datafusion/pull/22237 (strategy swap at row group 
boundaries). #22237 needed an arrow-rs `StrategySwap` API that did not ship. 
This PR uses `ParquetPushDecoder::into_builder` (arrow-rs 59.2+).
   
   ```mermaid
   graph LR
     A1["#25673 Optional wrapper"]
     A2["#25681 producers mark filters"]
     A3["#25674 gate"]
     A4["#25682 Parquet consumer"]
     B1["#25683 FilterExec consumer"]
     B2["#25721 FilterExec reordering"]
     C1["#25677 InList collapse"]
     C2["#25713 split join filter"]
     P["#22384 post-scan filter"]
     F["#25722 post-scan skips optional filters"]
     E2["#25726 per-conjunct pruning stats"]
     E3["E3 adaptive placement"]
     A1 --> A2
     C1 --> C2
     C2 --> A2
     A1 --> A4
     A3 --> A4
     A1 --> B1
     A3 --> B1
     A3 -. shares filter_stats .- B2
     P --> F
     A1 --> F
     A4 --> E3
     F --> E3
     E2 --> E3
     classDef this fill:#f6e7d6,stroke:#b25e12,stroke-width:3px
     class E3 this
   ```
   
   ## Rationale for this change
   
   With `pushdown_filters = true`, each pushable conjunct is a `RowFilter` 
predicate. This is not always faster:
   
   | Case | Cost of a row filter |
   |---|---|
   | The conjunct reads all output columns | Nothing to skip. A second decode 
pass or the predicate cache |
   | The conjunct removes rows in a scattered pattern | The decoder decodes and 
then filters. No skipped runs |
   | Remote object store | One more fetch round trip for each row filter stage 
and row group |
   | Optional filter that its gate paused (#25682) | The stage stays. Its 
columns are decoded and cached for nothing |
   
   Success criterion (for the benchmark bot): `pushdown_filters = true` with 
adaptive placement is never slower than `pushdown_filters = false`.
   
   ## What changes are included in this PR?
   
   New config `datafusion.execution.adaptive_filter_placement` (default 
`false`). With `pushdown_filters = true`, the scan decides for each conjunct, 
at file open and at each row group boundary:
   
   | Conjunct | Placements |
   |---|---|
   | Required | `RowFilter` or `PostScan` (the #22384 post-scan filter) |
   | Optional, `optional_filter_mode = adaptive` | `RowFilter` or `Skip` (while 
the gate is paused) |
   | Optional, `always` / `pruning_only` | Not changed (#25682) |
   | Rejected by the row filter | Required: `PostScan`. Optional: not used 
(#25722) |
   
   ```mermaid
   flowchart TD
       A[Conjunct at file open or row group boundary] --> B{Optional and gate 
paused?}
       B -- yes --> S[Skip]
       B -- no --> C{Optional?}
       C -- yes --> R[RowFilter]
       C -- no --> D{>= 8192 measured rows?}
       D -- no --> E{Stats pruned >= 50% of row groups<br/>or no unread output 
bytes?}
       E -- yes --> P[PostScan]
       E -- no --> R
       D -- yes --> F{benefit > cost, 10% margin}
       F -- yes --> R
       F -- no --> P
   ```
   
   ```text
   benefit (ns/row) = skippable fraction * unread output bytes per row * decode 
ns per byte
   cost    (ns/row) = mean fetch latency / rows of the next row group
   ```
   
   | Input | Source |
   |---|---|
   | Skippable fraction | Rows in 64-row windows where the conjunct passes no 
row. Measured in the row filter and in the post-scan filter, pooled over files 
and partitions |
   | Unread output bytes | Footer: output columns that the conjunct does not 
read |
   | Decode ns per byte | Measured decode time of the decoded batches |
   | Fetch latency | Measured `get_byte_ranges` time |
   | Statistics prior | #25726 `prune_with_conjunct_stats` during row group 
pruning (first production caller) |
   | Gate state | `OptionalFilterGate::is_paused` only. A skipped conjunct 
counts down the pause with `begin_batch()` for the batches of the skipped row 
groups |
   
   When the placement changes, the stream calls `into_builder()` at the row 
group boundary, sets the new `RowFilter` and projection mask, and builds a new 
`DecoderProjection`.
   
   | Commit | Content | Part |
   |---|---|---|
   | `feat: add filter_stats ...`, `feat: add OptionalFilterGate ...`, `feat: 
add optional_filter_mode config`, `feat: let the Parquet scan skip optional 
filters ...` | #25682 on top of #25722 | Base |
   | `test: end-to-end tests for optional filters in the Parquet post-scan 
path` | Tests that need both #25682 and #25722. Moves to the PR that merges 
second | Base |
   | 3 `pruning` commits | #25726 | Base |
   | `feat(parquet): adaptive placement of scan filter conjuncts` | Config, 
`filter_placement` module (`model`, `stats`, `FilePlacement`), stream and 
opener wiring, unit and opener tests | This PR |
   | `feat(parquet): seed filter placement with per-conjunct pruning 
statistics` | #25726 call in `RowGroupAccessPlanFilter` and the prior | This PR 
|
   | `test: sqllogictests for adaptive filter placement` | Same results with 
and without placement | This PR |
   
   | #22237 knob | This PR |
   |---|---|
   | `filter_pushdown_min_bytes_per_sec` | Removed. Measured benefit against 
measured fetch cost |
   | `filter_collecting_byte_ratio_threshold` | Removed. No unread output bytes 
means `PostScan` |
   | `filter_confidence_z` | Removed. 8192 minimum rows and a 10% margin |
   
   Limits: no change for a file with a live row selection (#24355), or when the 
new post-scan set changes the narrowed batch schema (nested columns). Optional 
conjuncts in `PostScan` are a follow-up.
   
   ## What is the testing strategy for this PR?
   
   | Test | Checks |
   |---|---|
   | `filter_placement::model` tests | Initial rule, statistics prior, measured 
benefit, fetch cost, margin (no flapping) |
   | `filter_placement::stats` tests | Skippable windows (nulls, partial 
window), pooled counts, keys by position or expression id |
   | `file_placement_follows_gate_and_measurements` | Paused gate gives `Skip`, 
pause countdown gives `RowFilter` again, scattered required conjunct moves to 
`PostScan` |
   | `scattered_filter_moves_post_scan_at_row_group_boundary` | Row filter for 
row group 0, post-scan for 1..3. `filter_placement_changes = 1`. Same rows as 
row filter only and post-scan only |
   | `clustered_filter_stays_row_filter` | No change |
   | `filter_on_all_output_columns_starts_post_scan` | Initial rule |
   | `statistics_prior_starts_clustered_filter_post_scan` | #25726 prior. No 
prior without statistics pruning |
   | `paused_optional_filter_is_removed_from_row_filter` | Fewer row filter 
rows than without placement. Same rows |
   | `limit_is_applied_after_placement_change` | Exact rows with `LIMIT` across 
a change |
   | `parquet_adaptive_filter_placement.slt` | Same results with and without 
placement |
   
   Benchmarks: run on the bot with `baseline: pushdown_filters=false` and 
`changed: pushdown_filters=true, adaptive_filter_placement=true`.
   
   ## Are there any user-facing changes?
   
   New config option (default `false`) and a new `filter_placement_changes` 
scan metric when it is on. No public API change.
   
   🤖 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]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to