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<FilterConjunct>"] 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]
