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

   > [!NOTE]
   > Prototype for discussion. It is an alternative to 
https://github.com/apache/datafusion/pull/25673 (the 
`OptionalFilterPhysicalExpr` wrapper). It does not close it.
   
   ## Which issue does this PR close?
   
   - Part of https://github.com/apache/datafusion/issues/22883 (optional 
filters).
   - It does not close an issue.
   
   ## Rationale for this change
   
   ### Motivation
   
   DataFusion passes every physical filter as one `Arc<dyn PhysicalExpr>`. An 
audit of `main` (`e8e41ae958`) found these problems:
   
   | # | Problem | Examples (`file:line` on `main`) | Fixed here? |
   |---|---|---|---|
   | 1 | A filter has no place for properties (optional, id, statistics). The 
only transport is the expression tree, so 
https://github.com/apache/datafusion/pull/25673 needs a wrapper expression that 
every rewriter, downcast and `.slt` plan must know. | producers: 
`hash_join/exec.rs:1987`, `sorts/sort.rs:1635`, `aggregates/mod.rs:2377` | yes: 
`FilterConjunct` |
   | 2 | Operators AND unaccepted filters together and lose their order and 
properties. | `physical-plan/src/filter.rs:799`, 
`datasource-parquet/src/source.rs:892`, `sorts/sort.rs:1677`, `union.rs:543` | 
yes |
   | 3 | Filters are built into one expression and split again: 11 
`split_conjunction` and 5 `conjunction` sites in physical code. | 
`filter.rs:728, 1049, 1127, 1183, 1531`, `row_filter.rs:495`, 
`page_filter.rs:166` | storage yes; opener: follow-up |
   | 4 | Equivalences and constants come from the whole predicate. A consumer 
that skips a conjunct would break them. | `filter.rs:460-503`, 
`file_scan_config/mod.rs:913-923` | `FilterExec` yes; scan: follow-up |
   | 5 | Nothing checks the rule that a node returns its parent filters in 
input order. | doc at `execution_plan.rs:845-850` | yes: debug check |
   | 6 | One conjunct that interval analysis rejects gives default selectivity 
to the whole predicate. | `filter.rs:398, 420` | follow-up |
   | 7 | Dynamic filters are found by tree walks, 2 to 3 times per Parquet 
file. | `opener/mod.rs:1854`, `push_decoder.rs:194`, `file_pruner.rs:97`, 
`execution_plan.rs:1130` | follow-up |
   | 8 | Per-conjunct metrics have no key. Row filter metrics are clones of one 
counter. | `row_filter.rs:538-559` | follow-up |
   
   The full audit (19 `PhysicalExpr` methods, 232 downcasts grouped by purpose, 
every split and re-join site) is below.
   
   <details>
   <summary>Full audit (<code>main</code> <code>e8e41ae958</code>)</summary>
   
   #### 2. `PhysicalExpr` methods that are about filters
   
   The trait (`physical-expr-common/src/physical_expr.rs:76`) has 19 methods. 6 
are mainly about filters (F).
   
   | Method | Line | Impls | Production callers | Better on a filter concept? |
   |---|---|---|---|---|
   | `evaluate_bounds` | 219 | 6 | interval graph (`cp_solver.rs:644`), 
`get_properties` (ordering) | no: also used for ordering of projections |
   | `propagate_constraints` | 249 | 6 | only `cp_solver.rs:692` (FilterExec 
`analyze`, symmetric hash join filter) | the node hook stays; the entry point 
(`analyze` + `check_support`) is filter-only |
   | `evaluate_statistics` / `propagate_statistics` | 280 / 333 | 4 / 2 | only 
`ExprStatisticsGraph`, which has no users | remove (dead, deprecated since 
54.0.0) |
   | `snapshot` | 432 | 2 | `snapshot_physical_expr_opt` (only 
`pruning_predicate.rs:531`) | partly: a dynamic filter can be nested in an 
expression |
   | `snapshot_generation` | 448 | 2 | only the deprecated 
`is_dynamic_physical_expr` (0 callers) | yes, or remove |
   | `expression_id` | 496 | 2 | dynamic filter identity 
(`execution_plan.rs:1138, 1621`, `aggregates/mod.rs:2402`, 
`hash_join/exec.rs:2024`, proto) | mostly: a conjunct id; proto dedup still 
needs expression identity |
   
   Conclusion: a filter concept does not remove these methods now. It gives a 
home for the per-conjunct parts (id, generation, analysis) in later PRs.
   
   #### 3. Downcasts that treat an expression as a filter
   
   232 production downcasts to built-in expression types in the physical 
crates. About 150 reason about a predicate.
   
   | Purpose | Sites | Files | Examples | Filter concept fixes it? |
   |---|---|---|---|---|
   | Find and track `DynamicFilterPhysicalExpr` | 12 | 10 | `tracker.rs:76`, 
`file_pruner.rs:97`, `opener/mod.rs:1854`, `push_decoder.rs:194`, 
`hash_join/exec.rs:2010, 2024` | yes: a per-conjunct `dynamic` property 
replaces the tree walks |
   | Split on AND (`split_conjunction`) | 13 | 6 | `filter.rs:728, 1049, 1127, 
1183, 1531`, `page_filter.rs:166`, `row_filter.rs:495`, 
`file_scan_config/mod.rs:1348` | yes: the conjuncts are stored split |
   | Re-join into one expression (loses properties) | 2 (+2 in 
FilterExec/Parquet) | 2 | `sort.rs:1677`, `union.rs:543` | yes |
   | FilterExec statistics (equality NDV, unique match, null-rejecting) | 14 | 
1 | `filter.rs:1139-1268` | yes: per-conjunct facts |
   | Interval support gating | 9 + 2 gates | 3 | `intervals/utils.rs:39-57`, 
gates at `filter.rs:398`, `sanity_checker.rs:132` | yes: per-conjunct analysis |
   | Equivalence and constants from a predicate | 10 | 3 | `filter.rs:530, 
548`, `utils/mod.rs:74, 81`, `file_scan_config/mod.rs:917, 1361` | yes: from 
required conjuncts only |
   | Pruning shapes (Column, Literal, InList, NOT, LIKE) | 30+ | 3 | 
`pruning_predicate.rs:847-2142`, `guarantee.rs:395-451` | no: expression 
semantics |
   | CASE-wrapped partitioned join filter | 0 consumers | — | built at 
`shared_bounds.rs:790`; pruning falls to the unhandled hook 
(`pruning_predicate.rs:1843`) | partly: per-partition conjunct data |
   | "No filter" as `lit(true)` | 6 | 5 | `filter.rs:852`, `topk/mod.rs:587`, 
`sort.rs:1294` | yes: an empty list |
   
   #### 4. Built into one expression, then split again
   
   | Site | What happens | Per-conjunct data that is lost | Struct fixes it? |
   |---|---|---|---|
   | `physical-optimizer/src/filter_pushdown.rs` (whole rule) | Carries 
`Vec<Arc<dyn PhysicalExpr>>` | properties of each filter | yes |
   | `physical-plan/src/filter.rs:728` | FilterExec splits its predicate to 
push self filters | properties | yes |
   | `physical-plan/src/filter.rs:792-799` | FilterExec re-joins unaccepted 
parent and self filters with `conjunction()` | order, properties | yes |
   | `physical-plan/src/filter.rs:366-503` | statistics and 
`compute_properties` split the predicate about 5 times per construction | 
cached per-conjunct facts | partly (home for a cache) |
   | `physical-plan/src/filter.rs:694-716` | projection swap rewrites the whole 
predicate; one bad conjunct blocks the swap | per-conjunct rewrite | partly |
   | `datasource-parquet/src/source.rs:890-896` | ParquetSource ANDs accepted 
filters with the old predicate | properties, "enforced" vs "pruning only" | yes 
|
   | `datasource-parquet/src/row_filter.rs:495`, `page_filter.rs:166` | opener 
splits the predicate again per file | properties, ids | follow-up (opener takes 
conjuncts) |
   | `pruning/src/pruning_predicate.rs:1847` | own AND recursion | which 
conjunct pruned | partly |
   | `pruning/src/file_pruner.rs:146`, `push_decoder.rs:229` | any dynamic 
change rebuilds the pruning predicate for all conjuncts | static vs dynamic 
conjuncts | follow-up |
   | `physical-plan/src/sorts/sort.rs:1677`, `union.rs:543` | wrap unaccepted 
filters in `FilterExec::try_new(conjunction(..))` | properties | yes |
   | `physical-optimizer/src/window_topn.rs:151` | needs the whole predicate to 
be `rn <= K`; `rn <= K AND x > 5` does not match | conjunct list | follow-up |
   
   Totals (non-test physical code): `split_conjunction` 11 sites, `conjunction` 
5 sites, `conjunction_opt` 1 site.
   
   #### 5. Other findings
   
   | Finding | Site | Note |
   |---|---|---|
   | Scan claims equivalences from a pruning-only predicate | 
`datasource/src/file_scan_config/mod.rs:913-923, 1348` | With 
`pushdown_filters=false` the scan does not enforce the predicate. Not verified 
to misplan. Follow-up. |
   | The pushdown protocol has no "retained" or "optional" state | 
`physical-plan/src/filter_pushdown.rs:96-146`, doc of 
`plan_contains_expression_id` (`execution_plan.rs:1130`) | Producers walk the 
plan to find if a consumer kept their filter. |
   | `reorder_filters` doc says "same order as written" | 
`common/src/config.rs:1345` | False after the re-joins at `filter.rs:792-799` 
and `source.rs:892`. |
   | Order of parent filters is a rule, but nothing checks it | 
`physical-plan/src/execution_plan.rs:845-850` (doc) | The optimizer only checks 
the length. |
   | `FilterRemapper` makes a new `Arc` for every column | 
`physical-plan/src/filter_pushdown.rs:426-445` | An unchanged filter loses its 
identity at every pass-through node. |
   
   </details>
   
   ### Idea
   
   Optionality is a property of a conjunct, not of an expression node. A 
`PhysicalFilter` is an ordered list of `FilterConjunct { expr, optional }`. It 
is not a `PhysicalExpr`, so rewriters, simplifiers and downcasts never see it.
   
   ```rust
   // datafusion_physical_expr::filter
   pub struct FilterConjunct { expr: Arc<dyn PhysicalExpr>, optional: bool }
   pub struct PhysicalFilter { conjuncts: Vec<FilterConjunct> }
   
   impl PhysicalFilter {
       pub fn conjuncts(&self) -> &[FilterConjunct];
       pub fn required(&self) -> impl Iterator<Item = &FilterConjunct>;
       pub fn optional(&self) -> impl Iterator<Item = &FilterConjunct>;
       pub fn to_expr(&self) -> Arc<dyn PhysicalExpr>;       // AND of all: 
evaluation, pruning, EXPLAIN
       pub fn required_expr(&self) -> Arc<dyn PhysicalExpr>; // AND of 
required: equivalence, constants
       pub fn try_map_exprs(self, f) -> Result<Self>;        // remap columns, 
keep properties
   }
   ```
   
   ### Transport: by position, no signature changes
   
   Nodes still see expressions in `gather_filters_for_pushdown`. The 
`FilterPushdown` rule holds `Vec<FilterConjunct>` and puts the properties back 
by position. Every node must already return one result per parent filter, in 
input order, or `PushedDown` is wrong today.
   
   ```mermaid
   flowchart LR
     P["producer<br/>with_optional_self_filter"] --> O["FilterPushdown 
rule<br/>Vec&lt;FilterConjunct&gt;"]
     O -- "expressions only" --> N["other nodes<br/>(no change)"]
     N -- "results in input order" --> O
     O --> D["DataSourceExec<br/>try_pushdown_filter(PhysicalFilter)"]
     D --> Q["ParquetSource<br/>stores PhysicalFilter"]
     D --> X["other sources<br/>default: flags dropped"]
     O --> F["FilterExec / SortExec / UnionExec<br/>keep the flags"]
   ```
   
   When a flag is lost (a source without the new method, an old proto reader, 
FFI), the filter stays and becomes required. That is correct, only slower.
   
   ## What changes are included in this PR?
   
   | Commit | Change |
   |---|---|
   | 1 | `physical_expr::filter::{PhysicalFilter, FilterConjunct}` |
   | 2 | `FilterPushdown` carries `FilterConjunct`s by position. 
`ChildFilterDescription::with_optional_self_filter`. 
`ChildFilterPushdownResult::new` / `conjunct()`. Debug check for the order 
rule. `FilterRemapper` keeps the same `Arc` for an unchanged filter, so the 
check sees pass-through nodes. |
   | 3 | `DataSource::try_pushdown_filter` and 
`FileSource::try_pushdown_filter`. The defaults drop the flags and call 
`try_pushdown_filters`. `FileScanConfig` remaps each conjunct. |
   | 4 | `FilterExec` stores a `PhysicalFilter`. Equivalences and constants use 
`required_expr()`. `FilterExec::try_new` and `predicate()` do not change. |
   | 5 | `SortExec` and `UnionExec` keep the flags when they add a 
`FilterExec`. |
   | 6 | `ParquetSource` stores a `PhysicalFilter` (no more `conjunction()` 
merge). The opener still gets `to_expr()`. |
   | 7 | Proto: `FilterExecNode.optional_conjuncts`, 
`ParquetScanExecNode.optional_predicate_conjuncts` (one flag per AND term). 
Additive: old readers apply all terms. |
   | 8 | Hash join (probe side), TopK and aggregate dynamic filters are 
optional. |
   | 9 | Tests. |
   
   No consumer skips an optional conjunct yet. Plans, EXPLAIN text and results 
do not change.
   
   ### Alternatives considered
   
   | Variant | Diff (`src` files incl. in-file unit tests, no generated code) | 
API break | Why not chosen |
   |---|---|---|---|
   | **This PR**: struct `FilterConjunct`, transport by position | 14 files, 
+940 / -172 | struct literal of `ChildFilterPushdownResult` (1 site in the 
repo) | chosen |
   | Trait `dyn PhysicalFilter { expr, is_optional, with_expr, evaluate }`, 
same transport (built and tested locally) | 14 files, +958 / -172 | same | 
Every consumer needs `expr()` and `with_expr()`, so the trait has one 
implementation: the struct behind `dyn`. Proto needs a codec for other 
implementations. Equality and `Debug` need `dyn` helpers. |
   | Trait with only `filter(&RecordBatch)` and `is_optional()` | not buildable 
| — | These consumers need the expression: optimizer remap and volatility check 
(`physical-optimizer/src/filter_pushdown.rs:463, 471, 535`), `FileScanConfig` 
unproject (`file_scan_config/mod.rs:1057`), Parquet pushdown check, pruning and 
row filter (`datasource-parquet/src/source.rs:396, 804, 820, 887`), 
`FilterExec` equivalence, split and projection swap (`filter.rs:207, 231, 350, 
752, 821`), proto (`filter.rs:981`). Line numbers are on this branch. |
   | Pushdown signatures carry the new type 
(`gather_filters_for_pushdown(Vec<FilterConjunct>)`, 
`FileSource::try_pushdown_filters(PhysicalFilter)`) | measured, not finished | 
all 14 `gather_filters_for_pushdown` impls, 5 `try_pushdown_filters` impls, 29 
`FilterDescription` builder sites in 15 files, all external impls | Changing 
only `gather_filters_for_pushdown` gives 25 errors in 12 files of 
`datafusion-physical-plan` alone. Transport by position gives the same result 
without this churn. |
   | Wrapper expression (https://github.com/apache/datafusion/pull/25673) | 10 
files | none | The flag is inside the expression tree: consumers must find it 
on the root AND chain, and rewriters, downcasts, the simplifier, pruning and 
equivalence must look through it. `Optional(...)` shows in 152 `.slt` plan 
lines. New public expression and proto node. |
   
   ### Follow-ups (not in this PR)
   
   | Follow-up | Audit row |
   |---|---|
   | Parquet opener takes conjuncts: no split in `row_filter.rs` / 
`page_filter.rs`, classify dynamic filters once | 3, 7 |
   | `FileScanConfig` equivalences from required (and enforced) conjuncts only 
| 4 |
   | Per-conjunct interval analysis and selectivity in `FilterExec` | 6 |
   | Conjunct id for metrics and statistics | 8 |
   | Consumers skip optional conjuncts (the rest of the stack for 
https://github.com/apache/datafusion/issues/22883) | 1 |
   
   ## What is the testing strategy for this PR?
   
   | Test | What it shows |
   |---|---|
   | `core/tests/parquet/optional_filters.rs` (3 tests, SQL on Parquet files) | 
Hash join, TopK and aggregate dynamic filters arrive at `ParquetSource` as 
optional conjuncts. User filters stay required. Results are correct. |
   | `test_hashjoin_optional_dynamic_filter_default_source` | A source that 
uses the default `try_pushdown_filter` gets the dynamic filter as a normal 
filter. The join result is correct. |
   | `test_reordered_parent_filters_are_rejected` | A node that reorders its 
parent filters gives an error in debug builds. |
   | `roundtrip_filter_optional_conjuncts`, 
`roundtrip_parquet_exec_optional_predicate_conjuncts` | Flags survive a proto 
round trip. A reader without the new fields applies all conjuncts as required. |
   | unit tests in `physical_expr::filter`, `FilterExec`, `ParquetSource` | 
flags kept on remap, split and absorb; equivalences only from required 
conjuncts; same predicate as the old entry point |
   
   Local runs: `cargo test` for `physical-expr`, `physical-plan`, 
`physical-optimizer`, `datasource`, `datasource-parquet` (lib), `datafusion` 
`core_integration` and `parquet_integration`, `datafusion-proto 
--all-features`; full `sqllogictests` with no plan changes; `cargo clippy 
--all-targets --workspace --features avro,integration-tests,extended_tests -- 
-D warnings`; `cargo doc` with `-D warnings` for the changed crates.
   
   ## Are there any user-facing changes?
   
   No query-visible change.
   
   API changes:
   
   | Item | Kind |
   |---|---|
   | `ChildFilterPushdownResult` has a private field: construct it with 
`ChildFilterPushdownResult::new` | breaking for struct literals and exhaustive 
patterns (1 test site in the repo) |
   | `physical_expr::filter` module; `DataSource::try_pushdown_filter`, 
`FileSource::try_pushdown_filter` (default methods); 
`ChildFilterDescription::with_optional_self_filter` / `with_self_conjunct`; 
`FilterDescription::self_conjuncts`; `FilterExecBuilder::new_with_filter` / 
`with_filter`; `FilterExec::filter`; `ParquetSource::physical_filter` | 
additive |
   | Proto fields `FilterExecNode.optional_conjuncts = 12`, 
`ParquetScanExecNode.optional_predicate_conjuncts = 8` | additive |
   
   🤖 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