Skip to content

feat: add OptionalFilterGate to pause optional filters that cost more than they save - #25674

Draft
adriangb wants to merge 6 commits into
apache:mainfrom
pydantic:optional-filter-gate
Draft

adriangb wants to merge 6 commits into
apache:mainfrom
pydantic:optional-filter-gate

Conversation

@adriangb

@adriangb adriangb commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

feat: add OptionalFilterGate to pause optional filters that cost more than they save

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
  A3 --> 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 A3 this
Loading

Rationale for this change

Consumers of optional filters need one shared rule for when to skip a filter. The rule compares the cost of the filter with the work that it saves:

Case Rows that pass Cost Result of the rule
Hash join filter, all keys match 100% any pause (no rows removed)
Star join filter on a small dimension (TPC-DS Q65, Q67) 20% 3 to 7 ns per row; the join probe costs 3.5 to 8 ns per row pause (cost > measured probe work)
TPC-DS Q50 store_sales hash lookup 7% 196 ms; the scan without it takes 72 ms pause (cost > saving)
TopK threshold (ClickBench Q23) about 0% low keep
Min/max bounds that remove 12% 88% about 1 ns per row keep (saves 2.4 ns per row)

What changes are included in this PR?

stateDiagram-v2
    [*] --> Evaluate
    Evaluate --> Paused: window of MIN_OBSERVED_ROWS rows; no rows removed, or cost > saving
    Evaluate --> Paused: another gate of the plan site published a pause (first window or probe only)
    Evaluate --> Evaluate: worth its cost, reset N
    Paused --> Evaluate: after N batches (probe), N doubles
    Paused --> Evaluate: dynamic filter changes, reset N, clear the shared pause
Loading

Decision at the end of each window. A window has at least sample_batches evaluated batches and at least MIN_OBSERVED_ROWS (8192) evaluated rows. All decisions (pause, keep, probe) need a full window. After a selective row filter a batch can have 2 to 7 rows, and the fixed cost of each call then looks like 600 to 8000 ns per row (ClickBench Q23).

Rule Pause when
No saving the window removed no rows
Cost cost > rows_removed × saving_ns_per_row (× 1.1 to pause, × 0.9 to resume)
Term Value
cost evaluation time + rows_in × overhead_ns_per_row
overhead_ns_per_row fixed cost for each evaluated row that the consumer measures (MeasuredRowSaving::set_overhead_ns_per_row, for example the cost of a Parquet row filter stage); 0 by default
saving_ns_per_row producer work + measured saving of the consumer
producer work the work that the producer of the filter does for each row that the filter removes (RemovedRowWork, for example the hash and the hash table lookup of a probe row). The smallest measured value of the dynamic filters in the filter. optional_filter_min_saving_ns_per_row is used only until the producer has measured MIN_OBSERVED_ROWS rows
measured saving of the consumer MeasuredRowSaving::set_ns_per_row, for example the Parquet decode time of the columns that the filter does not read

Shared pauses (SharedGateVerdict, one atomic word for each plan site):

Gate state Uses a pause that another gate published
new gate yes: it starts paused
first window, or probe window after a pause yes: it stops the window and pauses
gate that keeps the filter (evidence of its own) no: with skewed data the filter can be worth its cost for some files only
filter changed clears the shared pause (it was measured on the old filter) and starts again
Item Scope Purpose
filter_stats::{FilterCost, Clock} (first commit) shared with FilterExec reordering rows in, rows out, time; injectable clock (ManualClock in tests)
filter_stats::MIN_OBSERVED_ROWS shared with the Parquet filter placement (#25727) the smallest sample for an adaptive decision (8192 rows, one default batch)
filter_stats::RemovedRowWork one per DynamicFilterPhysicalExpr, shared by all derived filters (removed_row_work()) lock-free rows and nanoseconds that the producer records; ns_per_row() is None before MIN_OBSERVED_ROWS rows
OptionalFilterGate one per stream begin_batch(), record(rows_in, rows_out, elapsed), pauses(), is_paused(), with_shared_verdict()
SharedGateVerdict one per plan site (for example one optional filter of one scan) the last pause (or end of a pause) of the gates of the site, in one AtomicU64
MeasuredRowSaving one per consumer filter lock-free f64 saving and overhead that the consumer updates
DynamicFilterTracking::classify + watcher().changed() change detection one tree walk when the gate is created, then one atomic load per watched filter per batch
Config Default
datafusion.execution.optional_filter_min_saving_ns_per_row 20 (a hash probe measured about 13 ns per row). A prior: used only until the producer of the filter has a measurement

Changes since the previous version:

Removed Why
optional_filter_max_pass_ratio and the pass ratio check the cost rule covers it; cheap bounds that remove 12% now stay on
pooled statistics (OptionalFilterSiteStats, PackedVerdict) and the tracker generation tag kept on a branch for a later PR; only the shared pauses (SharedGateVerdict) are in this PR
optional_filter_mode moved to the consumer PRs (#25682, #25683)
test-only public API (evaluate, rows_skipped, batches_evaluated, next_pause_batches, site) no caller
Added Why
windows of at least MIN_OBSERVED_ROWS rows small batches after a selective row filter made cheap filters look expensive (ClickBench Q23)
SharedGateVerdict each file paid for its own first window and its own probes (TPC-H Q9: five join filters that remove no rows, in each of 12 files)
measured overhead for each evaluated row a Parquet row filter stage costs more than the evaluation of a cheap predicate
RemovedRowWork as the saving the fixed 20 ns kept filters on that cost more than the join work that they saved (TPC-DS Q65, Q67, Q90: 8-18% slower than pruning_only)

No consumer uses the gate in this PR.

What is the testing strategy for this PR?

Deterministic unit tests with a manual clock:

  • a filter that removes no rows pauses; a free filter that removes 1% stays on
  • backoff doubles up to the maximum and resets after a selective probe
  • a dynamic filter update restarts evaluation; a static or complete filter is never watched
  • a TopK-like filter that tightens over time stays on; skewed input is found by a probe
  • no backoff growth while paused (regression test)
  • cost rule: a selective but expensive filter pauses; cheap filters stay on (also the 12% case); a higher measured saving turns a paused filter back on; the 1.1 / 0.9 margins
  • windows: no decision before MIN_OBSERVED_ROWS rows (window_needs_min_observed_rows)
  • measured overhead adds to the cost (measured_overhead_adds_to_cost)
  • producer work: a measured RemovedRowWork replaces the configured saving (producer_work_replaces_configured_saving, removed_row_work_needs_min_observed_rows)
  • shared pauses: a new gate starts paused; a first window and a probe use a shared pause; a gate that keeps the filter ignores it; a filter change clears it; the packed word round-trips

Config docs and information_schema.slt are updated.

Are there any user-facing changes?

One new config option. No behavior change: no consumer uses the gate in this PR. cargo-semver-checks flags the new ExecutionOptions field, which is expected for a new config option.

🤖 Generated with Claude Code

@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 proto Related to proto crate labels Sep 24, 2026
@github-actions

github-actions Bot commented Sep 24, 2026 •

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 [  28.958s] (current)
     Parsing datafusion-common v55.1.0 (current)
      Parsed [   0.053s] (current)
    Building datafusion-common v55.1.0 (baseline)
       Built [  27.041s] (baseline)
     Parsing datafusion-common v55.1.0 (baseline)
      Parsed [   0.054s] (baseline)
    Checking datafusion-common v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.851s] 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.optional_filter_min_saving_ns_per_row 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 [  58.299s] datafusion-common
    Building datafusion-physical-expr v55.1.0 (current)
       Built [  23.936s] (current)
     Parsing datafusion-physical-expr v55.1.0 (current)
      Parsed [   0.045s] (current)
    Building datafusion-physical-expr v55.1.0 (baseline)
       Built [  23.662s] (baseline)
     Parsing datafusion-physical-expr v55.1.0 (baseline)
      Parsed [   0.043s] (baseline)
    Checking datafusion-physical-expr v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.411s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [  50.115s] datafusion-physical-expr
    Building datafusion-sqllogictest v55.1.0 (current)
       Built [  81.849s] (current)
     Parsing datafusion-sqllogictest v55.1.0 (current)
      Parsed [   0.019s] (current)
    Building datafusion-sqllogictest v55.1.0 (baseline)
       Built [  81.800s] (baseline)
     Parsing datafusion-sqllogictest v55.1.0 (baseline)
      Parsed [   0.020s] (baseline)
    Checking datafusion-sqllogictest v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.093s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [ 168.329s] datafusion-sqllogictest

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

codecov-commenter commented Sep 24, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 98.92473% with 9 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.52%. Comparing base (1be6b04) to head (87974dc).
⚠️ Report is 18 commits behind head on main.

Files with missing lines Patch % Lines
...tafusion/physical-expr/src/optional_filter_gate.rs 99.18% 4 Missing and 2 partials ⚠️
datafusion/physical-expr/src/filter_stats.rs 96.66% 3 Missing ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #25674      +/-   ##
==========================================
+ Coverage   82.49%   82.52%   +0.03%     
==========================================
  Files        1140     1142       +2     
  Lines      438781   439365     +584     
  Branches   438781   439365     +584     
==========================================
+ Hits       361965   362597     +632     
+ Misses      54938    54892      -46     
+ Partials    21878    21876       -2     

☔ 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 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>
… than they save

Add a runtime gate that pauses optional filters (filters that are not
needed for correctness, such as hash join and TopK dynamic filters) when
they cost more than they save. The gate is a per-stream state machine
(Evaluate / Paused with exponential backoff) that restarts evaluation
when the filter changes. The gate finds the dynamic filters one time
with `DynamicFilterTracking::classify` and then polls their
subscriptions, so a check does not walk the filter tree. Gates do not
share state.

At the end of each window of evaluated batches the gate pauses the
filter if the window removed no rows, or if its evaluation time is
larger than the work that the removed rows save:
`(rows_in - rows_out) * saving_ns_per_row`. The saving for each row is
the configured minimum plus an optional value that the consumer measures
and updates (`MeasuredRowSaving`). The cost rule has a margin (pause
above 1.1x the saving, resume below 0.9x) so that a filter does not
switch on and off when cost and saving are almost equal.

Add the `datafusion.execution.optional_filter_min_saving_ns_per_row`
option (default 20). No operator uses the gate yet, so behavior does not
change.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@adriangb
adriangb force-pushed the optional-filter-gate branch from 5f5cc7c to 5760ad5 Compare September 24, 2026 23:35
@adriangb adriangb changed the title feat: add OptionalFilterGate and optional_filter_mode config feat: add OptionalFilterGate to pause optional filters that cost more than they save Sep 24, 2026
@github-actions github-actions Bot removed the proto Related to proto crate label Sep 24, 2026
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 25, 2026
A consumer can now give the gate a fixed cost for each evaluated row in
addition to the evaluation time, with
`MeasuredRowSaving::set_overhead_ns_per_row`. The gate adds it to the cost
of each window. The Parquet scan uses it for the fixed cost of a row filter
stage, which is larger than the evaluation time of a cheap predicate.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 25, 2026
Each gate paid for its own first window and its own probes. A scan opens
its files at the same time, thus a filter that removes nothing cost one
window in each file (TPC-H Q9: five join filters that remove no rows and a
CASE routing filter that costs 83 ns for each row, in each of 12 files).

`SharedGateVerdict` holds the last pause (or end of a pause) of the gates
of one plan site in one atomic word. A gate without evidence of its own
(before its first decision, after a filter change, after a pause) uses a
pause that another gate published after the last verdict that it saw: a
new gate starts paused, and a gate in its first window or in a probe
window stops and pauses. Thus usually only one gate probes after a pause.
A gate that keeps the filter does not use the pauses of other gates
(skewed data). A filter change clears the shared pause, because it was
measured on the old filter.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 25, 2026
A gate decided after `sample_batches` batches, whatever their size. After
a selective row filter a batch can have 2 to 7 rows, and the fixed cost
of each call then looks like 600 to 8000 ns for each row: ClickBench Q23
paused the TopK filter on `EventTime` (0.4 ns for each row on full
batches) on such windows, and the shared verdict spread these pauses to
the other files.

All decisions (pause, keep, probe) now need a window of at least
`sample_batches` batches and `MIN_OBSERVED_ROWS` rows. The constant moves
to `filter_stats`, so that the gate and the Parquet filter placement use
the same sample size. All published shared pauses come from such windows.

PR: apache#25674 (the `FilterExec` test data change belongs to the FilterExec
optional filter commit of the same stack)

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 25, 2026
The gate assumed that each removed row saves `min_saving_ns_per_row`
(20 ns) after the filter. For hash join dynamic filters this is the probe
work of the join, and it is much smaller in star joins with small
dimension tables: 3.5 to 8 ns for each probe row on the TPC-DS SF1
`date_dim` joins (Q65, Q67), 2 ns on Q90, 17 ns on TPC-H Q9. Filters that
cost 3 to 7 ns for each row and remove 80% of the rows thus stayed on,
and cost more than the join work they saved (8-18% slower than
`pruning_only` on the bot).

`RemovedRowWork` (in `filter_stats`) is the work that the producer of a
filter does for each row that the filter removes, as the producer
measures it. Each `DynamicFilterPhysicalExpr` has one, shared by all its
derived filters. The gate uses the smallest measured work of the dynamic
filters in its filter as the saving of a removed row, and
`min_saving_ns_per_row` only until the producer has measured
`MIN_OBSERVED_ROWS` rows (a prior).

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb and others added 2 commits September 25, 2026 15:12
A consumer can now give the gate a fixed cost for each evaluated row in
addition to the evaluation time, with
`MeasuredRowSaving::set_overhead_ns_per_row`. The gate adds it to the cost
of each window. The Parquet scan uses it for the fixed cost of a row filter
stage, which is larger than the evaluation time of a cheap predicate.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Each gate paid for its own first window and its own probes. A scan opens
its files at the same time, thus a filter that removes nothing cost one
window in each file (TPC-H Q9: five join filters that remove no rows and a
CASE routing filter that costs 83 ns for each row, in each of 12 files).

`SharedGateVerdict` holds the last pause (or end of a pause) of the gates
of one plan site in one atomic word. A gate without evidence of its own
(before its first decision, after a filter change, after a pause) uses a
pause that another gate published after the last verdict that it saw: a
new gate starts paused, and a gate in its first window or in a probe
window stops and pauses. Thus usually only one gate probes after a pause.
A gate that keeps the filter does not use the pauses of other gates
(skewed data). A filter change clears the shared pause, because it was
measured on the old filter.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 25, 2026
A consumer can now give the gate a fixed cost for each evaluated row in
addition to the evaluation time, with
`MeasuredRowSaving::set_overhead_ns_per_row`. The gate adds it to the cost
of each window. The Parquet scan uses it for the fixed cost of a row filter
stage, which is larger than the evaluation time of a cheap predicate.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 25, 2026
Each gate paid for its own first window and its own probes. A scan opens
its files at the same time, thus a filter that removes nothing cost one
window in each file (TPC-H Q9: five join filters that remove no rows and a
CASE routing filter that costs 83 ns for each row, in each of 12 files).

`SharedGateVerdict` holds the last pause (or end of a pause) of the gates
of one plan site in one atomic word. A gate without evidence of its own
(before its first decision, after a filter change, after a pause) uses a
pause that another gate published after the last verdict that it saw: a
new gate starts paused, and a gate in its first window or in a probe
window stops and pauses. Thus usually only one gate probes after a pause.
A gate that keeps the filter does not use the pauses of other gates
(skewed data). A filter change clears the shared pause, because it was
measured on the old filter.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 25, 2026
A gate decided after `sample_batches` batches, whatever their size. After
a selective row filter a batch can have 2 to 7 rows, and the fixed cost
of each call then looks like 600 to 8000 ns for each row: ClickBench Q23
paused the TopK filter on `EventTime` (0.4 ns for each row on full
batches) on such windows, and the shared verdict spread these pauses to
the other files.

All decisions (pause, keep, probe) now need a window of at least
`sample_batches` batches and `MIN_OBSERVED_ROWS` rows. The constant moves
to `filter_stats`, so that the gate and the Parquet filter placement use
the same sample size. All published shared pauses come from such windows.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 25, 2026
The gate assumed that each removed row saves `min_saving_ns_per_row`
(20 ns) after the filter. For hash join dynamic filters this is the probe
work of the join, and it is much smaller in star joins with small
dimension tables: 3.5 to 8 ns for each probe row on the TPC-DS SF1
`date_dim` joins (Q65, Q67), 2 ns on Q90, 17 ns on TPC-H Q9. Filters that
cost 3 to 7 ns for each row and remove 80% of the rows thus stayed on,
and cost more than the join work they saved (8-18% slower than
`pruning_only` on the bot).

`RemovedRowWork` (in `filter_stats`) is the work that the producer of a
filter does for each row that the filter removes, as the producer
measures it. Each `DynamicFilterPhysicalExpr` has one, shared by all its
derived filters. The gate uses the smallest measured work of the dynamic
filters in its filter as the saving of a removed row, and
`min_saving_ns_per_row` only until the producer has measured
`MIN_OBSERVED_ROWS` rows (a prior).

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
A consumer can now give the gate a fixed cost for each evaluated row in
addition to the evaluation time, with
`MeasuredRowSaving::set_overhead_ns_per_row`. The gate adds it to the cost
of each window. The Parquet scan uses it for the fixed cost of a row filter
stage, which is larger than the evaluation time of a cheap predicate.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
Each gate paid for its own first window and its own probes. A scan opens
its files at the same time, thus a filter that removes nothing cost one
window in each file (TPC-H Q9: five join filters that remove no rows and a
CASE routing filter that costs 83 ns for each row, in each of 12 files).

`SharedGateVerdict` holds the last pause (or end of a pause) of the gates
of one plan site in one atomic word. A gate without evidence of its own
(before its first decision, after a filter change, after a pause) uses a
pause that another gate published after the last verdict that it saw: a
new gate starts paused, and a gate in its first window or in a probe
window stops and pauses. Thus usually only one gate probes after a pause.
A gate that keeps the filter does not use the pauses of other gates
(skewed data). A filter change clears the shared pause, because it was
measured on the old filter.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
A gate decided after `sample_batches` batches, whatever their size. After
a selective row filter a batch can have 2 to 7 rows, and the fixed cost
of each call then looks like 600 to 8000 ns for each row: ClickBench Q23
paused the TopK filter on `EventTime` (0.4 ns for each row on full
batches) on such windows, and the shared verdict spread these pauses to
the other files.

All decisions (pause, keep, probe) now need a window of at least
`sample_batches` batches and `MIN_OBSERVED_ROWS` rows. The constant moves
to `filter_stats`, so that the gate and the Parquet filter placement use
the same sample size. All published shared pauses come from such windows.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
The gate assumed that each removed row saves `min_saving_ns_per_row`
(20 ns) after the filter. For hash join dynamic filters this is the probe
work of the join, and it is much smaller in star joins with small
dimension tables: 3.5 to 8 ns for each probe row on the TPC-DS SF1
`date_dim` joins (Q65, Q67), 2 ns on Q90, 17 ns on TPC-H Q9. Filters that
cost 3 to 7 ns for each row and remove 80% of the rows thus stayed on,
and cost more than the join work they saved (8-18% slower than
`pruning_only` on the bot).

`RemovedRowWork` (in `filter_stats`) is the work that the producer of a
filter does for each row that the filter removes, as the producer
measures it. Each `DynamicFilterPhysicalExpr` has one, shared by all its
derived filters. The gate uses the smallest measured work of the dynamic
filters in its filter as the saving of a removed row, and
`min_saving_ns_per_row` only until the producer has measured
`MIN_OBSERVED_ROWS` rows (a prior).

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
A consumer can now give the gate a fixed cost for each evaluated row in
addition to the evaluation time, with
`MeasuredRowSaving::set_overhead_ns_per_row`. The gate adds it to the cost
of each window. The Parquet scan uses it for the fixed cost of a row filter
stage, which is larger than the evaluation time of a cheap predicate.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
Each gate paid for its own first window and its own probes. A scan opens
its files at the same time, thus a filter that removes nothing cost one
window in each file (TPC-H Q9: five join filters that remove no rows and a
CASE routing filter that costs 83 ns for each row, in each of 12 files).

`SharedGateVerdict` holds the last pause (or end of a pause) of the gates
of one plan site in one atomic word. A gate without evidence of its own
(before its first decision, after a filter change, after a pause) uses a
pause that another gate published after the last verdict that it saw: a
new gate starts paused, and a gate in its first window or in a probe
window stops and pauses. Thus usually only one gate probes after a pause.
A gate that keeps the filter does not use the pauses of other gates
(skewed data). A filter change clears the shared pause, because it was
measured on the old filter.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
A gate decided after `sample_batches` batches, whatever their size. After
a selective row filter a batch can have 2 to 7 rows, and the fixed cost
of each call then looks like 600 to 8000 ns for each row: ClickBench Q23
paused the TopK filter on `EventTime` (0.4 ns for each row on full
batches) on such windows, and the shared verdict spread these pauses to
the other files.

All decisions (pause, keep, probe) now need a window of at least
`sample_batches` batches and `MIN_OBSERVED_ROWS` rows. The constant moves
to `filter_stats`, so that the gate and the Parquet filter placement use
the same sample size. All published shared pauses come from such windows.

PR: apache#25674 (the `FilterExec` test data change belongs to the FilterExec
optional filter commit of the same stack)

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
The gate assumed that each removed row saves `min_saving_ns_per_row`
(20 ns) after the filter. For hash join dynamic filters this is the probe
work of the join, and it is much smaller in star joins with small
dimension tables: 3.5 to 8 ns for each probe row on the TPC-DS SF1
`date_dim` joins (Q65, Q67), 2 ns on Q90, 17 ns on TPC-H Q9. Filters that
cost 3 to 7 ns for each row and remove 80% of the rows thus stayed on,
and cost more than the join work they saved (8-18% slower than
`pruning_only` on the bot).

`RemovedRowWork` (in `filter_stats`) is the work that the producer of a
filter does for each row that the filter removes, as the producer
measures it. Each `DynamicFilterPhysicalExpr` has one, shared by all its
derived filters. The gate uses the smallest measured work of the dynamic
filters in its filter as the saving of a removed row, and
`min_saving_ns_per_row` only until the producer has measured
`MIN_OBSERVED_ROWS` rows (a prior).

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
A consumer can now give the gate a fixed cost for each evaluated row in
addition to the evaluation time, with
`MeasuredRowSaving::set_overhead_ns_per_row`. The gate adds it to the cost
of each window. The Parquet scan uses it for the fixed cost of a row filter
stage, which is larger than the evaluation time of a cheap predicate.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
Each gate paid for its own first window and its own probes. A scan opens
its files at the same time, thus a filter that removes nothing cost one
window in each file (TPC-H Q9: five join filters that remove no rows and a
CASE routing filter that costs 83 ns for each row, in each of 12 files).

`SharedGateVerdict` holds the last pause (or end of a pause) of the gates
of one plan site in one atomic word. A gate without evidence of its own
(before its first decision, after a filter change, after a pause) uses a
pause that another gate published after the last verdict that it saw: a
new gate starts paused, and a gate in its first window or in a probe
window stops and pauses. Thus usually only one gate probes after a pause.
A gate that keeps the filter does not use the pauses of other gates
(skewed data). A filter change clears the shared pause, because it was
measured on the old filter.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
A gate decided after `sample_batches` batches, whatever their size. After
a selective row filter a batch can have 2 to 7 rows, and the fixed cost
of each call then looks like 600 to 8000 ns for each row: ClickBench Q23
paused the TopK filter on `EventTime` (0.4 ns for each row on full
batches) on such windows, and the shared verdict spread these pauses to
the other files.

All decisions (pause, keep, probe) now need a window of at least
`sample_batches` batches and `MIN_OBSERVED_ROWS` rows. The constant moves
to `filter_stats`, so that the gate and the Parquet filter placement use
the same sample size. All published shared pauses come from such windows.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
The gate assumed that each removed row saves `min_saving_ns_per_row`
(20 ns) after the filter. For hash join dynamic filters this is the probe
work of the join, and it is much smaller in star joins with small
dimension tables: 3.5 to 8 ns for each probe row on the TPC-DS SF1
`date_dim` joins (Q65, Q67), 2 ns on Q90, 17 ns on TPC-H Q9. Filters that
cost 3 to 7 ns for each row and remove 80% of the rows thus stayed on,
and cost more than the join work they saved (8-18% slower than
`pruning_only` on the bot).

`RemovedRowWork` (in `filter_stats`) is the work that the producer of a
filter does for each row that the filter removes, as the producer
measures it. Each `DynamicFilterPhysicalExpr` has one, shared by all its
derived filters. The gate uses the smallest measured work of the dynamic
filters in its filter as the saving of a removed row, and
`min_saving_ns_per_row` only until the producer has measured
`MIN_OBSERVED_ROWS` rows (a prior).

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
A consumer can now give the gate a fixed cost for each evaluated row in
addition to the evaluation time, with
`MeasuredRowSaving::set_overhead_ns_per_row`. The gate adds it to the cost
of each window. The Parquet scan uses it for the fixed cost of a row filter
stage, which is larger than the evaluation time of a cheap predicate.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
Each gate paid for its own first window and its own probes. A scan opens
its files at the same time, thus a filter that removes nothing cost one
window in each file (TPC-H Q9: five join filters that remove no rows and a
CASE routing filter that costs 83 ns for each row, in each of 12 files).

`SharedGateVerdict` holds the last pause (or end of a pause) of the gates
of one plan site in one atomic word. A gate without evidence of its own
(before its first decision, after a filter change, after a pause) uses a
pause that another gate published after the last verdict that it saw: a
new gate starts paused, and a gate in its first window or in a probe
window stops and pauses. Thus usually only one gate probes after a pause.
A gate that keeps the filter does not use the pauses of other gates
(skewed data). A filter change clears the shared pause, because it was
measured on the old filter.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
A gate decided after `sample_batches` batches, whatever their size. After
a selective row filter a batch can have 2 to 7 rows, and the fixed cost
of each call then looks like 600 to 8000 ns for each row: ClickBench Q23
paused the TopK filter on `EventTime` (0.4 ns for each row on full
batches) on such windows, and the shared verdict spread these pauses to
the other files.

All decisions (pause, keep, probe) now need a window of at least
`sample_batches` batches and `MIN_OBSERVED_ROWS` rows. The constant moves
to `filter_stats`, so that the gate and the Parquet filter placement use
the same sample size. All published shared pauses come from such windows.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
The gate assumed that each removed row saves `min_saving_ns_per_row`
(20 ns) after the filter. For hash join dynamic filters this is the probe
work of the join, and it is much smaller in star joins with small
dimension tables: 3.5 to 8 ns for each probe row on the TPC-DS SF1
`date_dim` joins (Q65, Q67), 2 ns on Q90, 17 ns on TPC-H Q9. Filters that
cost 3 to 7 ns for each row and remove 80% of the rows thus stayed on,
and cost more than the join work they saved (8-18% slower than
`pruning_only` on the bot).

`RemovedRowWork` (in `filter_stats`) is the work that the producer of a
filter does for each row that the filter removes, as the producer
measures it. Each `DynamicFilterPhysicalExpr` has one, shared by all its
derived filters. The gate uses the smallest measured work of the dynamic
filters in its filter as the saving of a removed row, and
`min_saving_ns_per_row` only until the producer has measured
`MIN_OBSERVED_ROWS` rows (a prior).

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 27, 2026
A consumer can now give the gate a fixed cost for each evaluated row in
addition to the evaluation time, with
`MeasuredRowSaving::set_overhead_ns_per_row`. The gate adds it to the cost
of each window. The Parquet scan uses it for the fixed cost of a row filter
stage, which is larger than the evaluation time of a cheap predicate.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 27, 2026
Each gate paid for its own first window and its own probes. A scan opens
its files at the same time, thus a filter that removes nothing cost one
window in each file (TPC-H Q9: five join filters that remove no rows and a
CASE routing filter that costs 83 ns for each row, in each of 12 files).

`SharedGateVerdict` holds the last pause (or end of a pause) of the gates
of one plan site in one atomic word. A gate without evidence of its own
(before its first decision, after a filter change, after a pause) uses a
pause that another gate published after the last verdict that it saw: a
new gate starts paused, and a gate in its first window or in a probe
window stops and pauses. Thus usually only one gate probes after a pause.
A gate that keeps the filter does not use the pauses of other gates
(skewed data). A filter change clears the shared pause, because it was
measured on the old filter.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 27, 2026
A gate decided after `sample_batches` batches, whatever their size. After
a selective row filter a batch can have 2 to 7 rows, and the fixed cost
of each call then looks like 600 to 8000 ns for each row: ClickBench Q23
paused the TopK filter on `EventTime` (0.4 ns for each row on full
batches) on such windows, and the shared verdict spread these pauses to
the other files.

All decisions (pause, keep, probe) now need a window of at least
`sample_batches` batches and `MIN_OBSERVED_ROWS` rows. The constant moves
to `filter_stats`, so that the gate and the Parquet filter placement use
the same sample size. All published shared pauses come from such windows.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 27, 2026
The gate assumed that each removed row saves `min_saving_ns_per_row`
(20 ns) after the filter. For hash join dynamic filters this is the probe
work of the join, and it is much smaller in star joins with small
dimension tables: 3.5 to 8 ns for each probe row on the TPC-DS SF1
`date_dim` joins (Q65, Q67), 2 ns on Q90, 17 ns on TPC-H Q9. Filters that
cost 3 to 7 ns for each row and remove 80% of the rows thus stayed on,
and cost more than the join work they saved (8-18% slower than
`pruning_only` on the bot).

`RemovedRowWork` (in `filter_stats`) is the work that the producer of a
filter does for each row that the filter removes, as the producer
measures it. Each `DynamicFilterPhysicalExpr` has one, shared by all its
derived filters. The gate uses the smallest measured work of the dynamic
filters in its filter as the saving of a removed row, and
`min_saving_ns_per_row` only until the producer has measured
`MIN_OBSERVED_ROWS` rows (a prior).

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 28, 2026
A consumer can now give the gate a fixed cost for each evaluated row in
addition to the evaluation time, with
`MeasuredRowSaving::set_overhead_ns_per_row`. The gate adds it to the cost
of each window. The Parquet scan uses it for the fixed cost of a row filter
stage, which is larger than the evaluation time of a cheap predicate.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 28, 2026
Each gate paid for its own first window and its own probes. A scan opens
its files at the same time, thus a filter that removes nothing cost one
window in each file (TPC-H Q9: five join filters that remove no rows and a
CASE routing filter that costs 83 ns for each row, in each of 12 files).

`SharedGateVerdict` holds the last pause (or end of a pause) of the gates
of one plan site in one atomic word. A gate without evidence of its own
(before its first decision, after a filter change, after a pause) uses a
pause that another gate published after the last verdict that it saw: a
new gate starts paused, and a gate in its first window or in a probe
window stops and pauses. Thus usually only one gate probes after a pause.
A gate that keeps the filter does not use the pauses of other gates
(skewed data). A filter change clears the shared pause, because it was
measured on the old filter.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 28, 2026
A gate decided after `sample_batches` batches, whatever their size. After
a selective row filter a batch can have 2 to 7 rows, and the fixed cost
of each call then looks like 600 to 8000 ns for each row: ClickBench Q23
paused the TopK filter on `EventTime` (0.4 ns for each row on full
batches) on such windows, and the shared verdict spread these pauses to
the other files.

All decisions (pause, keep, probe) now need a window of at least
`sample_batches` batches and `MIN_OBSERVED_ROWS` rows. The constant moves
to `filter_stats`, so that the gate and the Parquet filter placement use
the same sample size. All published shared pauses come from such windows.

PR: apache#25674

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
adriangb added a commit to pydantic/datafusion that referenced this pull request Sep 28, 2026
The gate assumed that each removed row saves `min_saving_ns_per_row`
(20 ns) after the filter. For hash join dynamic filters this is the probe
work of the join, and it is much smaller in star joins with small
dimension tables: 3.5 to 8 ns for each probe row on the TPC-DS SF1
`date_dim` joins (Q65, Q67), 2 ns on Q90, 17 ns on TPC-H Q9. Filters that
cost 3 to 7 ns for each row and remove 80% of the rows thus stayed on,
and cost more than the join work they saved (8-18% slower than
`pruning_only` on the bot).

`RemovedRowWork` (in `filter_stats`) is the work that the producer of a
filter does for each row that the filter removes, as the producer
measures it. Each `DynamicFilterPhysicalExpr` has one, shared by all its
derived filters. The gate uses the smallest measured work of the dynamic
filters in its filter as the saving of a removed row, and
`min_saving_ns_per_row` only until the producer has measured
`MIN_OBSERVED_ROWS` rows (a prior).

PR: apache#25674

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 sqllogictest SQL Logic Tests (.slt)

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants