Skip to content

feat: adaptive conjunct reordering in FilterExec - #25721

Draft
adriangb wants to merge 2 commits into
apache:mainfrom
pydantic:filter-exec-reordering
Draft

adriangb wants to merge 2 commits into
apache:mainfrom
pydantic:filter-exec-reordering

Conversation

@adriangb

@adriangb adriangb commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

feat: adaptive conjunct reordering in FilterExec

Which issue does this PR close?

graph LR
  FS["filter_stats commit"]
  A1["#25673 Optional wrapper"]
  A2["#25681 producers mark filters"]
  A3["#25674 gate"]
  A4["#25682 Parquet consumer"]
  B1["#25683 FilterExec consumer"]
  B2["#25721 FilterExec reordering"]
  C2["#25713 split join filter"]
  P["#22384 post-scan filter"]
  F["#25722 post-scan skips optional filters"]
  FS --> A3
  FS --> B2
  A1 --> A2
  C2 --> A2
  A1 --> A4
  A3 --> A4
  A1 --> B1
  A3 --> B1
  A1 --> F
  P --> F
  classDef this fill:#f6e7d6,stroke:#b25e12,stroke-width:3px
  class B2 this
Loading

Rationale for this change

The order of the conjuncts of an AND predicate is important. BinaryExpr evaluates the right side only on the rows that the left side keeps, if the left side keeps few rows. A selective, cheap conjunct written last makes the other conjuncts run on all rows.

-- regexp_like(s, 'rare') removes most rows, but it is written last
SELECT * FROM t WHERE regexp_like(s, 'a') AND regexp_like(s, 'rare');

This PR starts from the review of #22698:

#22698 review point This PR
One AND evaluator, not a second engine the result is a plain BinaryExpr AND chain
The cost model must match BinaryExpr pre-selection uses PRE_SELECTION_THRESHOLD (now pub, #[doc(hidden)])
State must not leak across executions the state is per stream
Deterministic tests injectable Clock (ManualClock in tests)
No unsafe unwrap, no UInt32Array downcast removed

What changes are included in this PR?

datafusion.execution.adaptive_filter_reordering (experimental, default false):

graph LR
  W["8 warm-up batches:<br/>measure each conjunct"] --> R["rank by rows removed per ns"]
  R --> D{"estimated cost ≥ 5% lower?"}
  D -- yes --> N["new BinaryExpr AND chain"]
  D -- no --> O["keep the original predicate"]
Loading
Item Purpose
filter_stats::{FilterCost, Clock, SystemClock, ManualClock} (first commit) rows in, rows out, time; injectable clock
filter/conjunct_order.rs ConjunctOrder: warm-up measurement, ranking, cost estimate with the BinaryExpr pre-selection rule
FilterExec uses ConjunctOrder for each stream when the option is on
Metric adaptive_reorders number of streams that changed the order

Volatile predicates are never reordered. A fallible conjunct can see different rows after a reorder (for example b <> 0 AND 1 / b > 2); the config documentation says so.

There are no optional-filter concepts in this PR.

What is the testing strategy for this PR?

Area Tests
conjunct_order.rs 9 scenario tests with a manual clock: a selective conjunct moves first, the best written order is kept, a cheap conjunct stays before an expensive selective one, no reorder when AND cannot pre-select, nulls prevent a reorder, right-nested result, empty batches are not measured, volatile predicates are not reordered, the pre-selection model matches BinaryExpr
FilterExec adaptive_filter_reordering_returns_same_rows (same rows, metric present only with two or more conjuncts)
adaptive_filter_reordering.slt same results with the option on and off
filter_stats cost derivations, manual clock, system clock

Are there any user-facing changes?

One new config option (default off) and one new metric. The default behavior does not change.

🤖 Generated with Claude Code

adriangb and others added 2 commits September 24, 2026 14:52
Add the `datafusion_physical_expr::filter_stats` module with the shared
primitives that adaptive filter code uses to measure filters at runtime:

- `Clock`: a monotonic clock in nanoseconds that tests can replace.
  `SystemClock` is the real clock. `ManualClock` moves only when a test
  moves it, thus decisions that use time are deterministic in tests.
- `FilterCost`: the rows in, the rows out and the evaluation time of one
  filter, and the derived cost for each row and rows removed for each
  nanosecond.
- `duration_nanos`: a `Duration` in nanoseconds, saturated to `u64::MAX`.

No code uses the module yet, thus behavior does not change.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Add adaptive conjunct reordering behind
`datafusion.execution.adaptive_filter_reordering` (default `false`): each
stream measures its conjuncts over a warm-up, ranks them by rows dropped
per nanosecond, and adopts a new order (built as a plain `BinaryExpr` AND
chain) only if the estimated cost, using `BinaryExpr`'s pre-selection
rule, is at least 5% lower. Predicates with volatile expressions are
never reordered. The decision is made one time for each stream, and
streams do not share state.

The measurements use `Clock` and `FilterCost` of
`datafusion_physical_expr::filter_stats`, so tests use a `ManualClock`.
`PRE_SELECTION_THRESHOLD` is exported (doc-hidden) so that the cost model
uses the same rule as `BinaryExpr`.

New metric: `adaptive_reorders`.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@github-actions github-actions Bot added documentation Improvements or additions to documentation physical-expr Changes to the physical-expr crates sqllogictest SQL Logic Tests (.slt) common Related to common crate physical-plan Changes to the physical-plan crate labels Sep 24, 2026
@github-actions

Copy link
Copy Markdown

Thank you for opening this pull request!

Reviewer note: cargo-semver-checks reported the current version number is not SemVer-compatible with the changes in this pull request (compared against the base branch).

Details
     Cloning apache/main
    Building datafusion-common v55.1.0 (current)
       Built [  33.677s] (current)
     Parsing datafusion-common v55.1.0 (current)
      Parsed [   0.068s] (current)
    Building datafusion-common v55.1.0 (baseline)
       Built [  32.908s] (baseline)
     Parsing datafusion-common v55.1.0 (baseline)
      Parsed [   0.068s] (baseline)
    Checking datafusion-common v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   1.096s] 223 checks: 222 pass, 1 fail, 0 warn, 31 skip

--- failure constructible_struct_adds_field: struct exhaustively constructible through public API adds field ---

Description:
A pub struct that could be exhaustively constructed with a literal using only public API has a new pub field, breaking existing exhaustive literals.
        ref: https://doc.rust-lang.org/reference/expressions/struct-expr.html
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.50.0/src/lints/constructible_struct_adds_field.ron

Failed in:
  field ExecutionOptions.adaptive_filter_reordering in /home/runner/work/datafusion/datafusion/datafusion/common/src/config.rs:896

     Summary semver requires new major version: 1 major and 0 minor checks failed
    Finished [  69.400s] datafusion-common
    Building datafusion-physical-expr v55.1.0 (current)
       Built [  28.569s] (current)
     Parsing datafusion-physical-expr v55.1.0 (current)
      Parsed [   0.054s] (current)
    Building datafusion-physical-expr v55.1.0 (baseline)
       Built [  28.643s] (baseline)
     Parsing datafusion-physical-expr v55.1.0 (baseline)
      Parsed [   0.053s] (baseline)
    Checking datafusion-physical-expr v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.519s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [  58.840s] datafusion-physical-expr
    Building datafusion-physical-plan v55.1.0 (current)
       Built [  37.854s] (current)
     Parsing datafusion-physical-plan v55.1.0 (current)
      Parsed [   0.179s] (current)
    Building datafusion-physical-plan v55.1.0 (baseline)
       Built [  37.940s] (baseline)
     Parsing datafusion-physical-plan v55.1.0 (baseline)
      Parsed [   0.181s] (baseline)
    Checking datafusion-physical-plan v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   1.057s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [  78.583s] datafusion-physical-plan
    Building datafusion-sqllogictest v55.1.0 (current)
       Built [  99.333s] (current)
     Parsing datafusion-sqllogictest v55.1.0 (current)
      Parsed [   0.023s] (current)
    Building datafusion-sqllogictest v55.1.0 (baseline)
       Built [  99.655s] (baseline)
     Parsing datafusion-sqllogictest v55.1.0 (baseline)
      Parsed [   0.025s] (baseline)
    Checking datafusion-sqllogictest v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.120s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [ 202.017s] datafusion-sqllogictest

@github-actions github-actions Bot added the auto detected api change Auto detected API change label Sep 24, 2026
@codecov-commenter

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 88.35490% with 63 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.49%. Comparing base (1be6b04) to head (319fcb7).
⚠️ Report is 7 commits behind head on main.

Files with missing lines Patch % Lines
...afusion/physical-plan/src/filter/conjunct_order.rs 88.05% 44 Missing and 4 partials ⚠️
datafusion/physical-plan/src/filter.rs 82.60% 4 Missing and 8 partials ⚠️
datafusion/physical-expr/src/filter_stats.rs 95.71% 3 Missing ⚠️
Additional details and impacted files
@@           Coverage Diff            @@
##             main   #25721    +/-   ##
========================================
  Coverage   82.49%   82.49%            
========================================
  Files        1140     1142     +2     
  Lines      438781   439571   +790     
  Branches   438781   439571   +790     
========================================
+ Hits       361965   362633   +668     
- Misses      54938    55019    +81     
- Partials    21878    21919    +41     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
The post-scan filter of the Parquet scan uses the pre-selection threshold
of `AND` to decide when to compact its working batch (next commit). This
is the same change as in the FilterExec reordering PR (apache#25721), which
models the same rule. When one of the two PRs merges, this commit becomes
empty in the other.

PR: apache#25727

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

auto detected api change Auto detected API change common Related to common crate documentation Improvements or additions to documentation physical-expr Changes to the physical-expr crates physical-plan Changes to the physical-plan crate sqllogictest SQL Logic Tests (.slt)

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants