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]
