Skip to content

Reapply "Add ExecutionPlan::apply_expressions() (apache#20337)" (apache#22437) - #24018

Open
jayshrivastava wants to merge 3 commits into
apache:mainfrom
jayshrivastava:js/dynamic-filters-method
Open

Reapply "Add ExecutionPlan::apply_expressions() (apache#20337)" (apache#22437)#24018
jayshrivastava wants to merge 3 commits into
apache:mainfrom
jayshrivastava:js/dynamic-filters-method

Conversation

@jayshrivastava

@jayshrivastava jayshrivastava commented Jul 30, 2026

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Rationale for this change

See the issue above.

What changes are included in this PR?

There's 3 commits in this PR:

Firstly, commit 1 re-applies the changes in #20337 (reverted in #22437).

Some of the reasons for why the original PR was reverted include
(a) apply_expressions is too complicated to implement and there's no concrete need to justify this complexity
(b) there was no usage of apply_expressions inside this repo

To address (a), justification for this method is provided in #23814. Additionally, commit 2 in this PR updates the apply_expressions API to be more ergonomic:

  • apply_expressions now yields &Arc<dyn PhysicalExpr> instead of &dyn PhysicalExpr. It should be easier for callers to use the new API without worrying about the lifetime restriction on a non Arc reference
  • adds a helper apply_expression_roots which makes abstracts away the TreeNodeRecursion complexity from implementors. Now, it's very trivial to implement ex.
fn apply_expressions(
    &self,
    f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> Result<TreeNodeRecursion>,
) -> Result<TreeNodeRecursion> {
    apply_expression_roots([&self.predicate_1], f)
    apply_expression_roots([&self.predicate_2], f)
    apply_expression_roots([&self.other_expression], f)
}

To address (b), commit 3 adds a usage of apply_expressions in physical-plan/src/aggregates/mod.rs. Previously, there was a hack that checked if a filter was pushed down using Arc::strong_count(dyn_filter) > 1. Now it uses apply_expressions.

Are these changes tested?

Yes.

Are there any user-facing changes?

There's a new mandatory method ExecutionPlan::apply_expressions(). See the upgrading guide and documentation for details.

@github-actions github-actions Bot added documentation Improvements or additions to documentation optimizer Optimizer rules core Core DataFusion crate catalog Related to the catalog crate proto Related to proto crate datasource Changes to the datasource crate ffi Changes to the ffi crate physical-plan Changes to the physical-plan crate labels Jul 30, 2026
@jayshrivastava jayshrivastava changed the title Js/dynamic filters method Reapply "Add ExecutionPlan::apply_expressions() (apache#20337)" (apache#22437) Jul 30, 2026
@github-actions

github-actions Bot commented Jul 30, 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 v54.1.0 (current)
       Built [ 106.204s] (current)
     Parsing datafusion v54.1.0 (current)
      Parsed [   0.035s] (current)
    Building datafusion v54.1.0 (baseline)
       Built [ 103.746s] (baseline)
     Parsing datafusion v54.1.0 (baseline)
      Parsed [   0.033s] (baseline)
    Checking datafusion v54.1.0 -> v54.1.0 (no change; assume patch)
     Checked [   0.616s] 223 checks: 223 pass, 30 skip
     Summary no semver update required
    Finished [ 212.635s] datafusion
    Building datafusion-catalog v54.1.0 (current)
       Built [  40.260s] (current)
     Parsing datafusion-catalog v54.1.0 (current)
      Parsed [   0.024s] (current)
    Building datafusion-catalog v54.1.0 (baseline)
       Built [  40.195s] (baseline)
     Parsing datafusion-catalog v54.1.0 (baseline)
      Parsed [   0.024s] (baseline)
    Checking datafusion-catalog v54.1.0 -> v54.1.0 (no change; assume patch)
     Checked [   0.112s] 223 checks: 223 pass, 30 skip
     Summary no semver update required
    Finished [  81.877s] datafusion-catalog
    Building datafusion-datasource v54.1.0 (current)
       Built [  40.523s] (current)
     Parsing datafusion-datasource v54.1.0 (current)
      Parsed [   0.031s] (current)
    Building datafusion-datasource v54.1.0 (baseline)
       Built [  40.526s] (baseline)
     Parsing datafusion-datasource v54.1.0 (baseline)
      Parsed [   0.033s] (baseline)
    Checking datafusion-datasource v54.1.0 -> v54.1.0 (no change; assume patch)
     Checked [   0.244s] 223 checks: 222 pass, 1 fail, 0 warn, 30 skip

--- failure trait_method_added: pub trait method added ---

Description:
A non-sealed public trait added a new method without a default implementation, which breaks downstream implementations of the trait
        ref: https://doc.rust-lang.org/cargo/reference/semver.html#trait-new-item-no-default
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.49.0/src/lints/trait_method_added.ron

Failed in:
  trait method datafusion_datasource::source::DataSource::apply_expressions in file /home/runner/work/datafusion/datafusion/datafusion/datasource/src/source.rs:240
  trait method datafusion_datasource::file::FileSource::apply_expressions in file /home/runner/work/datafusion/datafusion/datafusion/datasource/src/file.rs:372

     Summary semver requires new major version: 1 major and 0 minor checks failed
    Finished [  82.854s] datafusion-datasource
    Building datafusion-datasource-arrow v54.1.0 (current)
       Built [  39.966s] (current)
     Parsing datafusion-datasource-arrow v54.1.0 (current)
      Parsed [   0.011s] (current)
    Building datafusion-datasource-arrow v54.1.0 (baseline)
       Built [  39.636s] (baseline)
     Parsing datafusion-datasource-arrow v54.1.0 (baseline)
      Parsed [   0.011s] (baseline)
    Checking datafusion-datasource-arrow v54.1.0 -> v54.1.0 (no change; assume patch)
     Checked [   0.071s] 223 checks: 223 pass, 30 skip
     Summary no semver update required
    Finished [  81.287s] datafusion-datasource-arrow
    Building datafusion-datasource-avro v54.1.0 (current)
       Built [  41.147s] (current)
     Parsing datafusion-datasource-avro v54.1.0 (current)
      Parsed [   0.010s] (current)
    Building datafusion-datasource-avro v54.1.0 (baseline)
       Built [  41.006s] (baseline)
     Parsing datafusion-datasource-avro v54.1.0 (baseline)
      Parsed [   0.010s] (baseline)
    Checking datafusion-datasource-avro v54.1.0 -> v54.1.0 (no change; assume patch)
     Checked [   0.081s] 223 checks: 223 pass, 30 skip
     Summary no semver update required
    Finished [  83.491s] datafusion-datasource-avro
    Building datafusion-datasource-csv v54.1.0 (current)
       Built [  40.388s] (current)
     Parsing datafusion-datasource-csv v54.1.0 (current)
      Parsed [   0.011s] (current)
    Building datafusion-datasource-csv v54.1.0 (baseline)
       Built [  39.911s] (baseline)
     Parsing datafusion-datasource-csv v54.1.0 (baseline)
      Parsed [   0.011s] (baseline)
    Checking datafusion-datasource-csv v54.1.0 -> v54.1.0 (no change; assume patch)
     Checked [   0.103s] 223 checks: 223 pass, 30 skip
     Summary no semver update required
    Finished [  81.565s] datafusion-datasource-csv
    Building datafusion-datasource-json v54.1.0 (current)
       Built [  40.022s] (current)
     Parsing datafusion-datasource-json v54.1.0 (current)
      Parsed [   0.012s] (current)
    Building datafusion-datasource-json v54.1.0 (baseline)
       Built [  39.850s] (baseline)
     Parsing datafusion-datasource-json v54.1.0 (baseline)
      Parsed [   0.013s] (baseline)
    Checking datafusion-datasource-json v54.1.0 -> v54.1.0 (no change; assume patch)
     Checked [   0.095s] 223 checks: 223 pass, 30 skip
     Summary no semver update required
    Finished [  81.349s] datafusion-datasource-json
    Building datafusion-datasource-parquet v54.1.0 (current)
       Built [  45.291s] (current)
     Parsing datafusion-datasource-parquet v54.1.0 (current)
      Parsed [   0.030s] (current)
    Building datafusion-datasource-parquet v54.1.0 (baseline)
       Built [  44.958s] (baseline)
     Parsing datafusion-datasource-parquet v54.1.0 (baseline)
      Parsed [   0.032s] (baseline)
    Checking datafusion-datasource-parquet v54.1.0 -> v54.1.0 (no change; assume patch)
     Checked [   0.159s] 223 checks: 223 pass, 30 skip
     Summary no semver update required
    Finished [  91.763s] datafusion-datasource-parquet
    Building datafusion-ffi v54.1.0 (current)
       Built [  61.595s] (current)
     Parsing datafusion-ffi v54.1.0 (current)
      Parsed [   0.061s] (current)
    Building datafusion-ffi v54.1.0 (baseline)
       Built [  61.122s] (baseline)
     Parsing datafusion-ffi v54.1.0 (baseline)
      Parsed [   0.061s] (baseline)
    Checking datafusion-ffi v54.1.0 -> v54.1.0 (no change; assume patch)
     Checked [   0.227s] 223 checks: 221 pass, 1 fail, 1 warn, 30 skip

--- failure constructible_struct_adds_field: externally-constructible struct adds field ---

Description:
A pub struct constructible with a struct literal has a new pub field. Existing struct literals must be updated to include the new field.
        ref: https://doc.rust-lang.org/reference/expressions/struct-expr.html
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.49.0/src/lints/constructible_struct_adds_field.ron

Failed in:
  field FFI_PhysicalExpr.expression_id in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/physical_expr/mod.rs:123
  field ForeignLibraryModule.create_exec_with_expressions in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/tests/mod.rs:112
  field FFI_ExecutionPlan.apply_expressions in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/execution_plan.rs:56
  field FFI_ExecutionPlan.version in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/execution_plan.rs:101

--- warning repr_c_plain_struct_fields_reordered: struct fields reordered in repr(C) struct ---

Description:
A public repr(C) struct had its fields reordered. This can change the struct's memory layout, possibly breaking FFI use cases that depend on field position and order.
        ref: https://doc.rust-lang.org/reference/type-layout.html#reprc-structs
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.49.0/src/lints/repr_c_plain_struct_fields_reordered.ron

Failed in:
  FFI_ExecutionPlan.with_new_children moved from position 3 to 4, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/execution_plan.rs:59
  FFI_ExecutionPlan.name moved from position 4 to 5, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/execution_plan.rs:63
  FFI_ExecutionPlan.execute moved from position 5 to 6, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/execution_plan.rs:67
  FFI_ExecutionPlan.repartitioned moved from position 6 to 7, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/execution_plan.rs:73
  FFI_ExecutionPlan.metrics moved from position 7 to 8, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/execution_plan.rs:82
  FFI_ExecutionPlan.partition_statistics moved from position 8 to 9, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/execution_plan.rs:88
  FFI_ExecutionPlan.clone moved from position 9 to 10, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/execution_plan.rs:95
  FFI_ExecutionPlan.release moved from position 10 to 11, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/execution_plan.rs:98
  FFI_ExecutionPlan.private_data moved from position 11 to 13, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/execution_plan.rs:105
  FFI_ExecutionPlan.library_marker_id moved from position 12 to 14, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/execution_plan.rs:110
  FFI_PhysicalExpr.display moved from position 17 to 18, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/physical_expr/mod.rs:126
  FFI_PhysicalExpr.hash moved from position 18 to 19, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/physical_expr/mod.rs:129
  FFI_PhysicalExpr.clone moved from position 19 to 20, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/physical_expr/mod.rs:133
  FFI_PhysicalExpr.release moved from position 20 to 21, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/physical_expr/mod.rs:136
  FFI_PhysicalExpr.version moved from position 21 to 22, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/physical_expr/mod.rs:139
  FFI_PhysicalExpr.private_data moved from position 22 to 23, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/physical_expr/mod.rs:143
  FFI_PhysicalExpr.library_marker_id moved from position 23 to 24, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/physical_expr/mod.rs:147
  ForeignLibraryModule.create_exec_with_statistics moved from position 15 to 16, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/tests/mod.rs:114
  ForeignLibraryModule.create_table_with_statistics moved from position 16 to 17, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/tests/mod.rs:116
  ForeignLibraryModule.create_physical_optimizer_rule moved from position 17 to 18, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/tests/mod.rs:119
  ForeignLibraryModule.create_context_aware_optimizer_rule moved from position 18 to 19, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/tests/mod.rs:121
  ForeignLibraryModule.version moved from position 19 to 20, in /home/runner/work/datafusion/datafusion/datafusion/ffi/src/tests/mod.rs:123

     Summary semver requires new major version: 1 major and 0 minor checks failed
     Warning produced 1 major and 0 minor level warnings
    Finished [ 124.722s] datafusion-ffi
    Building datafusion-physical-optimizer v54.1.0 (current)
       Built [  40.641s] (current)
     Parsing datafusion-physical-optimizer v54.1.0 (current)
      Parsed [   0.020s] (current)
    Building datafusion-physical-optimizer v54.1.0 (baseline)
       Built [  40.214s] (baseline)
     Parsing datafusion-physical-optimizer v54.1.0 (baseline)
      Parsed [   0.023s] (baseline)
    Checking datafusion-physical-optimizer v54.1.0 -> v54.1.0 (no change; assume patch)
     Checked [   0.126s] 223 checks: 223 pass, 30 skip
     Summary no semver update required
    Finished [  82.235s] datafusion-physical-optimizer
    Building datafusion-physical-plan v54.1.0 (current)
       Built [  38.034s] (current)
     Parsing datafusion-physical-plan v54.1.0 (current)
      Parsed [   0.141s] (current)
    Building datafusion-physical-plan v54.1.0 (baseline)
       Built [  38.268s] (baseline)
     Parsing datafusion-physical-plan v54.1.0 (baseline)
      Parsed [   0.143s] (baseline)
    Checking datafusion-physical-plan v54.1.0 -> v54.1.0 (no change; assume patch)
     Checked [   0.604s] 223 checks: 222 pass, 1 fail, 0 warn, 30 skip

--- failure trait_method_added: pub trait method added ---

Description:
A non-sealed public trait added a new method without a default implementation, which breaks downstream implementations of the trait
        ref: https://doc.rust-lang.org/cargo/reference/semver.html#trait-new-item-no-default
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.49.0/src/lints/trait_method_added.ron

Failed in:
  trait method datafusion_physical_plan::execution_plan::ExecutionPlan::apply_expressions in file /home/runner/work/datafusion/datafusion/datafusion/physical-plan/src/execution_plan.rs:331
  trait method datafusion_physical_plan::ExecutionPlan::apply_expressions in file /home/runner/work/datafusion/datafusion/datafusion/physical-plan/src/execution_plan.rs:331

     Summary semver requires new major version: 1 major and 0 minor checks failed
    Finished [  78.593s] datafusion-physical-plan
    Building datafusion-proto v54.1.0 (current)
       Built [  60.676s] (current)
     Parsing datafusion-proto v54.1.0 (current)
      Parsed [   0.017s] (current)
    Building datafusion-proto v54.1.0 (baseline)
       Built [  60.818s] (baseline)
     Parsing datafusion-proto v54.1.0 (baseline)
      Parsed [   0.018s] (baseline)
    Checking datafusion-proto v54.1.0 -> v54.1.0 (no change; assume patch)
     Checked [   0.256s] 223 checks: 223 pass, 30 skip
     Summary no semver update required
    Finished [ 123.298s] datafusion-proto

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

codecov-commenter commented Jul 30, 2026

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 44.01330% with 505 lines in your changes missing coverage. Please review.
✅ Project coverage is 80.76%. Comparing base (541caab) to head (5514194).
⚠️ Report is 5 commits behind head on main.

Files with missing lines Patch % Lines
datafusion/physical-plan/src/execution_plan.rs 57.39% 66 Missing and 6 partials ⚠️
datafusion/physical-plan/src/test/exec.rs 0.00% 36 Missing ⚠️
datafusion/core/src/physical_planner.rs 0.00% 24 Missing ⚠️
...sion/physical-optimizer/src/output_requirements.rs 0.00% 22 Missing ⚠️
datafusion/physical-plan/src/limit.rs 0.00% 22 Missing ⚠️
datafusion/physical-plan/src/sorts/sort.rs 0.00% 21 Missing ⚠️
datafusion/physical-plan/src/aggregates/mod.rs 76.47% 11 Missing and 9 partials ⚠️
datafusion/physical-plan/src/display.rs 0.00% 16 Missing ⚠️
...n/physical-plan/src/sorts/sort_preserving_merge.rs 0.00% 14 Missing ⚠️
...ysical-plan/src/windows/bounded_window_agg_exec.rs 0.00% 14 Missing ⚠️
... and 31 more
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #24018      +/-   ##
==========================================
- Coverage   80.84%   80.76%   -0.08%     
==========================================
  Files        1096     1099       +3     
  Lines      373936   375205    +1269     
  Branches   373936   375205    +1269     
==========================================
+ Hits       302313   303041     +728     
- Misses      53584    54084     +500     
- Partials    18039    18080      +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.

@jayshrivastava
jayshrivastava force-pushed the js/dynamic-filters-method branch from 8eff543 to b8794d2 Compare July 30, 2026 22:12
@jayshrivastava
jayshrivastava force-pushed the js/dynamic-filters-method branch from b8794d2 to 5514194 Compare July 30, 2026 23:26
/// `apply_expressions` + `downcast_ref::<DynamicFilterPhysicalExpr>` and
/// counting nodes. Neither API is observable from SQL.
#[tokio::test]
async fn test_discover_dynamic_filters_via_expressions_api() {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This test explicitly covers that dynamic filters are discoverable post-pushdown. This is what we want in datafusion-distributed and I assume other projects want this as well

aggregate.with_dynamic_filter_expr(dynamic_filter)?
} else {
let mut aggregate = aggregate;
aggregate.dynamic_filter = None;

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Semi-related. I think this was a bug. We kept the default dynamic filter created by try_new around even there's no pushed down filter. Now, we explicitly remove it. The test I added in datafusion/proto/tests/cases/roundtrip_physical_plan.rs covers this.

.flatten()
.map(|sort_expr| &sort_expr.expr),
f,
)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Question for reviewers - this pattern repeats a lot. Is it worth adding a helper?

apply_expression_roots(
    properties
        .output_ordering()
        .into_iter()
        .flatten()
        .map(|sort_expr| &sort_expr.expr),
    f
)

The projection pattern also repeats a bit

      crate::apply_expression_roots(
          self.projector
              .projection()
              .as_ref()
              .iter()
              .map(|proj_expr| &proj_expr.expr),
          f,
      )

@jayshrivastava

Copy link
Copy Markdown
Contributor Author

I see the code coverage is isn't very high. I don't see any good candidates to call apply_expressions generally. Maybe we can implement check_invariants on a few ExecutionPlan implementations which assert apply_expressions returns the right number of expressions.

@jayshrivastava
jayshrivastava marked this pull request as ready for review July 31, 2026 00:19
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 catalog Related to the catalog crate core Core DataFusion crate datasource Changes to the datasource crate documentation Improvements or additions to documentation ffi Changes to the ffi crate optimizer Optimizer rules physical-plan Changes to the physical-plan crate proto Related to proto crate

Projects

None yet

Development

Successfully merging this pull request may close these issues.

make dynamic filters in ExecutionPlan nodes discoverable

2 participants