From 06c503f917761d0e88ab9e5cedf13fa244d623a4 Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Thu, 17 Sep 2026 20:53:04 -0500 Subject: [PATCH 1/4] refactor: share `is_pure_extraction_projection` between the pushdown rules Take the expression list instead of a `LogicalPlan`, and make the function `pub(crate)`. `PushDownFilter` needs the same predicate in the next commit. No behaviour change. Co-Authored-By: Claude Fable 5.1 --- .../optimizer/src/extract_leaf_expressions.rs | 27 +++++++++---------- 1 file changed, 12 insertions(+), 15 deletions(-) diff --git a/datafusion/optimizer/src/extract_leaf_expressions.rs b/datafusion/optimizer/src/extract_leaf_expressions.rs index 6daec5a79db5e..fe07990d57067 100644 --- a/datafusion/optimizer/src/extract_leaf_expressions.rs +++ b/datafusion/optimizer/src/extract_leaf_expressions.rs @@ -1151,15 +1151,14 @@ fn split_and_push_projection( } } -/// Returns true if the plan is a Projection where ALL expressions are either -/// `Alias(EXTRACTED_EXPR_PREFIX, ...)` or `Column`, with at least one extraction. +/// Returns true if `exprs` are the expressions of a *pure extraction +/// projection*: every expression is either `Alias(EXTRACTED_EXPR_PREFIX, ...)` +/// or a bare `Column`, and there is at least one extraction alias. +/// /// Such projections can safely be pushed further without re-extraction. -fn is_pure_extraction_projection(plan: &LogicalPlan) -> bool { - let LogicalPlan::Projection(proj) = plan else { - return false; - }; +pub(crate) fn is_pure_extraction_projection(exprs: &[Expr]) -> bool { let mut has_extraction = false; - for expr in &proj.expr { + for expr in exprs { match expr { Expr::Alias(alias) if alias.name.starts_with(EXTRACTED_EXPR_PREFIX) => { has_extraction = true; @@ -1196,6 +1195,7 @@ fn push_extraction_pairs( proj_input, target_schema.as_ref(), )?; + let merged_is_pure = is_pure_extraction_projection(&merged.expr); let merged_plan = LogicalPlan::Projection(merged); // After merging, try to push the result further down, but ONLY @@ -1206,7 +1206,7 @@ fn push_extraction_pairs( // the (None, true) fallback can't find the original aliases. // This handles: Extraction → Recovery(cols) → Filter → ... → TableScan // by pushing through the recovery projection AND the filter in one pass. - if is_pure_extraction_projection(&merged_plan) + if merged_is_pure && let Some(pushed) = try_push_input(&merged_plan, alias_generator)? { return Ok(Some(pushed)); @@ -3657,14 +3657,11 @@ mod tests { fn test_is_pure_extraction_projection() -> Result<()> { let scan = test_table_scan_with_struct()?; - // Not a projection. - assert!(!is_pure_extraction_projection(&scan)); - // Columns only: nothing was extracted. let columns_only = LogicalPlanBuilder::from(scan.clone()) .project(vec![col("id"), col("user")])? .build()?; - assert!(!is_pure_extraction_projection(&columns_only)); + assert!(!is_pure_extraction_projection(&columns_only.expressions())); // One extraction alias and one pass-through column. let pure = LogicalPlanBuilder::from(scan.clone()) @@ -3673,7 +3670,7 @@ mod tests { col("id"), ])? .build()?; - assert!(is_pure_extraction_projection(&pure)); + assert!(is_pure_extraction_projection(&pure.expressions())); // The same shape under a `CommonSubexprEliminate` alias is not an // extraction projection. @@ -3683,13 +3680,13 @@ mod tests { col("id"), ])? .build()?; - assert!(!is_pure_extraction_projection(&common_expr)); + assert!(!is_pure_extraction_projection(&common_expr.expressions())); // A bare expression is never allowed. let bare_expr = LogicalPlanBuilder::from(scan) .project(vec![leaf_udf(col("user"), "status"), col("id")])? .build()?; - assert!(!is_pure_extraction_projection(&bare_expr)); + assert!(!is_pure_extraction_projection(&bare_expr.expressions())); Ok(()) } From 81c4d0537bfb0d7aac907d2b2f09e74bd8fa2dc0 Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Thu, 17 Sep 2026 20:57:31 -0500 Subject: [PATCH 2/4] fix: `PushDownFilter` yields to pure extraction projections `PushDownFilter` and `PushDownLeafProjections` want the opposite order for an adjacent filter and pure extraction projection, so on `main` they undo each other on every optimizer pass and the rule that runs later decides the plan. Give the extraction projection precedence: `rewrite_projection` keeps every predicate above a pure extraction projection. The projection then stays next to the scan, so a Parquet scan merges it into the file projection and reads only the struct leaf. This is the property the leaf pushdown feature exists for, and it is lost with the opposite precedence whenever `datafusion.execution.parquet.pushdown_filters` is `false`, which is the default. The filter loses nothing. `PushDownFilter` runs before `ExtractLeafExpressions`, so it records the predicate in `TableScan::filters` in the first pass, before any extraction projection exists. Record the invariant in the module documentation of both rules and in the query optimizer guide. Two expected plans in `projection_pushdown.slt` change: - A `TableScan` loses a `Boolean(true)` entry from `partial_filters`. The entry was a no-op. - `CAST(character_length(...) AS Int64)` prints as `CAST(character_length(...) AS length(get_field(...)) AS Int64)`. The SQL planner writes `length(x)` as `character_length(x) AS length(x)`, and the coercion pass wraps that alias in the cast. `main` prints the same text for `SELECT a * 2 + length(b) AS score FROM t`, with no struct and no extraction. The alias disappeared from this test only because of the projection merge that the rule fight caused. The physical plan does not change. Co-Authored-By: Claude Fable 5.1 --- .../optimizer/src/extract_leaf_expressions.rs | 39 ++++++++++++++++ datafusion/optimizer/src/push_down_filter.rs | 44 +++++++++++++++++++ .../test_files/projection_pushdown.slt | 13 +++--- .../library-user-guide/query-optimizer.md | 10 +++++ 4 files changed, 99 insertions(+), 7 deletions(-) diff --git a/datafusion/optimizer/src/extract_leaf_expressions.rs b/datafusion/optimizer/src/extract_leaf_expressions.rs index fe07990d57067..53d6663240fbe 100644 --- a/datafusion/optimizer/src/extract_leaf_expressions.rs +++ b/datafusion/optimizer/src/extract_leaf_expressions.rs @@ -19,6 +19,41 @@ //! access `user['status']`) closer to data sources, enabling early data reduction //! and source-level optimizations (e.g., Parquet column pruning). See //! [`ExtractLeafExpressions`] (pass 1) and [`PushDownLeafProjections`] (pass 2). +//! +//! # Precedence over [`PushDownFilter`] +//! +//! [`PushDownLeafProjections`] and [`PushDownFilter`] both move nodes towards +//! the leaves, and for an adjacent filter and *pure extraction projection* (a +//! projection whose expressions are only `__datafusion_extracted_N` aliases and +//! pass-through columns) they want the opposite order: +//! +//! ```text +//! Filter: t.date = '2025-01-03' <-- (A) +//! Projection: get_field(t.ids, 'id1') AS __datafusion_extracted_1, t.date <-- (B) +//! TableScan: t +//! ``` +//! +//! `PushDownFilter` wants (A) below (B). `PushDownLeafProjections` wants (B) +//! below (A). Before was +//! decided the two rules undid each other on every optimizer pass, and the rule +//! that runs later in the list won. +//! +//! **Invariant: pure extraction projections win.** `PushDownFilter` does not +//! move a filter below a pure extraction projection, so the plan above is the +//! final plan. `PushDownLeafProjections` keeps moving such a projection through +//! a filter, which is how the projection reaches the scan. +//! +//! The reason is that the extraction projection is the node a source absorbs. +//! A Parquet scan merges it into the file projection and reads only the struct +//! leaf. With the opposite precedence the scan reads the whole struct whenever +//! `datafusion.execution.parquet.pushdown_filters` is `false`, which is the +//! default. The filter loses nothing by staying one node higher, because +//! `PushDownFilter` records the predicate in +//! [`TableScan::filters`](datafusion_expr::logical_plan::TableScan) in the pass +//! that runs before the extraction projection exists, so row group pruning and +//! source level filtering still happen. +//! +//! [`PushDownFilter`]: crate::push_down_filter::PushDownFilter use indexmap::{IndexMap, IndexSet}; use std::collections::{BTreeSet, HashMap}; @@ -1156,6 +1191,10 @@ fn split_and_push_projection( /// or a bare `Column`, and there is at least one extraction alias. /// /// Such projections can safely be pushed further without re-extraction. +/// +/// [`PushDownFilter`](crate::push_down_filter::PushDownFilter) uses the same +/// predicate to keep a filter above such a projection. See the module +/// documentation for the precedence rule. pub(crate) fn is_pure_extraction_projection(exprs: &[Expr]) -> bool { let mut has_extraction = false; for expr in exprs { diff --git a/datafusion/optimizer/src/push_down_filter.rs b/datafusion/optimizer/src/push_down_filter.rs index b358c8abd2650..703307d6ee908 100644 --- a/datafusion/optimizer/src/push_down_filter.rs +++ b/datafusion/optimizer/src/push_down_filter.rs @@ -16,6 +16,28 @@ // under the License. //! [`PushDownFilter`] applies filters as early as possible +//! +//! # Precedence: pure extraction projections win +//! +//! [`PushDownLeafProjections`] moves a *pure extraction projection* towards the +//! leaves. Such a projection has only `__datafusion_extracted_N` aliases and +//! pass-through columns. [`PushDownFilter`] moves filters towards the leaves +//! too. For an adjacent filter and pure extraction projection the two rules +//! want the opposite order, so they undo each other on every optimizer pass. +//! +//! **Invariant: `PushDownFilter` yields to a pure extraction projection.** A +//! filter is never moved below such a projection. The extraction projection +//! stays at the bottom of the plan, next to the scan, and the filter stays +//! above it. +//! +//! The reason is that the extraction projection is the node a source absorbs. +//! A Parquet scan merges it into the file projection and reads only the struct +//! leaf, which is what the rule exists for. Keeping the filter one node higher +//! costs nothing at the scan, because `PushDownFilter` records the predicate in +//! [`TableScan::filters`](datafusion_expr::logical_plan::TableScan) in the pass +//! that runs before the extraction projection exists. +//! +//! [`PushDownLeafProjections`]: crate::extract_leaf_expressions::PushDownLeafProjections use std::collections::{HashMap, HashSet}; use std::sync::Arc; @@ -42,6 +64,7 @@ use datafusion_expr::{ TableProviderFilterPushDown, and, or, }; +use crate::extract_leaf_expressions::is_pure_extraction_projection; use crate::optimizer::ApplyOrder; use crate::simplify_expressions::{reorder_predicates, simplify_predicates}; use crate::utils::{ @@ -1388,6 +1411,27 @@ fn rewrite_projection( predicates: Vec, mut projection: Projection, ) -> Result<(Transformed, Vec)> { + // Precedence rule: a filter never moves below a pure extraction projection. + // + // `PushDownLeafProjections` moves such a projection below an adjacent + // filter, so a filter that moved below it is put back above it in the same + // optimizer pass. The two rules then undo each other on every pass until + // the pass limit stops them, and the surviving plan is decided by rule + // order alone. `PushDownFilter` yields here, because the extraction + // projection is the node the source absorbs: leaving it at the bottom keeps + // Parquet struct field pruning, and a filter kept one node higher still + // reaches the scan through `TableScan::filters`, which the pass that + // created the extraction projection has already set. + // + // See the module documentation of `extract_leaf_expressions` for the + // full statement of the invariant. + if is_pure_extraction_projection(&projection.expr) { + return Ok(( + Transformed::no(LogicalPlan::Projection(projection)), + predicates, + )); + } + // Partition projection expressions into non-pushable vs pushable. // Non-pushable expressions are volatile (must not be duplicated) or // MoveTowardsLeafNodes (cheap expressions like get_field where re-inlining diff --git a/datafusion/sqllogictest/test_files/projection_pushdown.slt b/datafusion/sqllogictest/test_files/projection_pushdown.slt index f0b8b9eeae73b..0c51b474c7cc5 100644 --- a/datafusion/sqllogictest/test_files/projection_pushdown.slt +++ b/datafusion/sqllogictest/test_files/projection_pushdown.slt @@ -1059,7 +1059,7 @@ query TT EXPLAIN SELECT s['value'] * 2 + length(s['label']) as score FROM simple_struct WHERE id > 1; ---- logical_plan -01)Projection: __datafusion_extracted_1 * Int64(2) + CAST(character_length(__datafusion_extracted_2) AS Int64) AS score +01)Projection: __datafusion_extracted_1 * Int64(2) + CAST(character_length(__datafusion_extracted_2) AS length(get_field(simple_struct.s, Utf8("label"))) AS Int64) AS score 02)--Filter: simple_struct.id > Int64(1) 03)----Projection: get_field(simple_struct.s, Utf8("value")) AS __datafusion_extracted_1, get_field(simple_struct.s, Utf8("label")) AS __datafusion_extracted_2, simple_struct.id 04)------TableScan: simple_struct projection=[id, s], partial_filters=[simple_struct.id > Int64(1)] @@ -2331,12 +2331,11 @@ logical_plan 05)--------SubqueryAlias: inner_t 06)----------Projection: simple_struct.id, simple_struct.s, __datafusion_extracted_1 07)------------Limit: skip=0, fetch=1 -08)--------------Filter: __datafusion_extracted_2 > Int64(120) AND __datafusion_extracted_1 != Utf8("delta") -09)----------------Filter: simple_struct.id = outer_ref(outer_t.id) -10)------------------Projection: get_field(simple_struct.s, Utf8("value")) AS __datafusion_extracted_2, get_field(simple_struct.s, Utf8("label")) AS __datafusion_extracted_1, simple_struct.id, simple_struct.s -11)--------------------TableScan: simple_struct, partial_filters=[simple_struct.id = outer_ref(outer_t.id), get_field(simple_struct.s, Utf8("value")) > Int64(120)] -12)--SubqueryAlias: outer_t -13)----TableScan: simple_struct projection=[id] +08)--------------Filter: simple_struct.id = outer_ref(outer_t.id) AND __datafusion_extracted_2 > Int64(120) AND __datafusion_extracted_1 != Utf8("delta") +09)----------------Projection: get_field(simple_struct.s, Utf8("value")) AS __datafusion_extracted_2, simple_struct.id, simple_struct.s, get_field(simple_struct.s, Utf8("label")) AS __datafusion_extracted_1 +10)------------------TableScan: simple_struct, partial_filters=[simple_struct.id = outer_ref(outer_t.id), get_field(simple_struct.s, Utf8("value")) > Int64(120)] +11)--SubqueryAlias: outer_t +12)----TableScan: simple_struct projection=[id] statement ok set datafusion.explain.logical_plan_only = false; diff --git a/docs/source/library-user-guide/query-optimizer.md b/docs/source/library-user-guide/query-optimizer.md index a3853d39ddced..883f922c06df9 100644 --- a/docs/source/library-user-guide/query-optimizer.md +++ b/docs/source/library-user-guide/query-optimizer.md @@ -162,6 +162,16 @@ individual operators or expressions. Sometimes there is an initial pass that visits the plan and builds state that is used in a second pass that performs the actual optimization. This approach is used in projection push down and filter push down. +### Rule Precedence + +Two rules can want the opposite order for the same pair of adjacent plan nodes. Each rule then undoes the work of the other one on every optimizer pass. The loop stops at `datafusion.optimizer.max_passes`, and the position of the two rules in the rule list decides the plan. This wastes a pass and makes the plan depend on the rule list. + +Do not let two rules compete. Give one rule precedence, make that rule yield, and record the decision in the module documentation of both rules. + +There is one such decision today: + +- `PushDownFilter` yields to a _pure extraction projection_. A pure extraction projection is a projection whose expressions are only `__datafusion_extracted_N` aliases and pass-through columns. `ExtractLeafExpressions` creates it, and `PushDownLeafProjections` moves it towards the leaves. `PushDownFilter` does not move a filter below such a projection, so the projection stays next to the scan. A Parquet scan then merges the projection into the file projection and reads only the struct leaf. The filter loses nothing, because `PushDownFilter` records the predicate in `TableScan::filters` in the pass that runs before the extraction projection exists. See [issue #14540](https://github.com/apache/datafusion/issues/14540). + ### Expression Naming Every expression in DataFusion has a name, which is used as the column name. For example, in this example the output From f146f1d752d32d2c85194d233d3a3305f36aaeee Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Thu, 17 Sep 2026 20:57:43 -0500 Subject: [PATCH 3/4] test: cover the filter and extraction projection precedence Unit tests in `push_down_filter.rs`: a filter stays above a pure extraction projection, and a projection that also computes an expression is still pushed through. A new section in `projection_pushdown.slt` with the query from https://github.com/apache/datafusion/issues/14540 on a memory table and on a Parquet file. The Parquet plan shows both properties the precedence must keep: `DataSourceExec` reads only the struct leaf `ids.id1`, and the `date` predicate reaches the scan for row group pruning. The section also pins the plan at a reduced `datafusion.optimizer.max_passes`. The simple shape gives the same plan at one pass as at the default. The issue shape gives the same plan at two passes, because the two filters merge in the second pass. Neither plan depends on the pass limit, so neither depends on which of the two rules runs last. Co-Authored-By: Claude Fable 5.1 --- datafusion/optimizer/src/push_down_filter.rs | 64 +++++ .../test_files/projection_pushdown.slt | 221 ++++++++++++++++++ 2 files changed, 285 insertions(+) diff --git a/datafusion/optimizer/src/push_down_filter.rs b/datafusion/optimizer/src/push_down_filter.rs index 703307d6ee908..370466f8ee971 100644 --- a/datafusion/optimizer/src/push_down_filter.rs +++ b/datafusion/optimizer/src/push_down_filter.rs @@ -4590,6 +4590,70 @@ mod tests { ) } + /// A filter is not moved below a pure extraction projection, even when its + /// predicate only references pass-through columns. + /// + /// `PushDownLeafProjections` moves such a projection back below the filter, + /// so pushing here would make the two rules undo each other on every + /// optimizer pass. See + /// and the module documentation. + #[test] + fn filter_not_pushed_through_pure_extraction_projection() -> Result<()> { + let table_scan = test_table_scan()?; + + // The shape `ExtractLeafExpressions` produces: one extraction alias + // plus pass-through columns. + let proj = LogicalPlanBuilder::from(table_scan) + .project(vec![ + leaf_udf_expr(col("a")).alias("__datafusion_extracted_1"), + col("b"), + col("c"), + ])? + .build()?; + + // `b` is a plain pass-through column, so without the precedence rule + // this predicate would reach the scan as a `full_filters` entry. + let plan = LogicalPlanBuilder::from(proj) + .filter(col("b").gt(lit(5i64)))? + .build()?; + + assert_optimized_plan_equal!( + plan, + @r" + Filter: test.b > Int64(5) + Projection: leaf_udf(test.a) AS __datafusion_extracted_1, test.b, test.c + TableScan: test + " + ) + } + + /// A projection that mixes an extraction alias with a computed expression is + /// not a pure extraction projection, so the filter still moves below it. + #[test] + fn filter_pushed_through_mixed_extraction_projection() -> Result<()> { + let table_scan = test_table_scan()?; + + let proj = LogicalPlanBuilder::from(table_scan) + .project(vec![ + leaf_udf_expr(col("a")).alias("__datafusion_extracted_1"), + (col("b") + lit(1i64)).alias("b_plus"), + col("c"), + ])? + .build()?; + + let plan = LogicalPlanBuilder::from(proj) + .filter(col("c").gt(lit(5i64)))? + .build()?; + + assert_optimized_plan_equal!( + plan, + @r" + Projection: leaf_udf(test.a) AS __datafusion_extracted_1, test.b + Int64(1) AS b_plus, test.c + TableScan: test, full_filters=[test.c > Int64(5)] + " + ) + } + #[test] fn filter_not_pushed_down_through_table_scan_with_fetch() -> Result<()> { let scan = test_table_scan()?; diff --git a/datafusion/sqllogictest/test_files/projection_pushdown.slt b/datafusion/sqllogictest/test_files/projection_pushdown.slt index 0c51b474c7cc5..2d5cb50af5798 100644 --- a/datafusion/sqllogictest/test_files/projection_pushdown.slt +++ b/datafusion/sqllogictest/test_files/projection_pushdown.slt @@ -2339,3 +2339,224 @@ logical_plan statement ok set datafusion.explain.logical_plan_only = false; + +##################### +# Section 17: Precedence between PushDownFilter and PushDownLeafProjections +# +# Both rules move nodes towards the leaves, and for an adjacent filter and +# "pure extraction projection" they want the opposite order. The precedence is: +# the extraction projection wins. It stays next to the scan, so the scan can +# absorb it and read only the struct leaf, and the filter stays above it. +# +# See https://github.com/apache/datafusion/issues/14540 +##################### + +statement ok +SET datafusion.execution.target_partitions = 1; + +statement ok +CREATE TABLE events_mem ( + "date" DATE, + "timestamp" TIMESTAMP, + ids STRUCT, + structs STRUCT +) AS VALUES + (DATE '2025-01-03', TIMESTAMP '2025-01-03 01:00:00', {id1: 'dev1', extra: 1}, {var1: 'user1', extra: 'e1'}), + (DATE '2025-01-03', TIMESTAMP '2025-01-03 02:00:00', {id1: 'dev1', extra: 2}, {var1: 'user2', extra: 'e2'}), + (DATE '2025-01-04', TIMESTAMP '2025-01-04 01:00:00', {id1: 'dev2', extra: 3}, {var1: 'user3', extra: 'e3'}); + +# The query from #14540. The extraction projection is the bottom node and the +# filter sits above it. +query TT +EXPLAIN WITH events AS ( + SELECT ids.id1 AS device, structs.var1 AS user, "timestamp" + FROM events_mem + WHERE "date" = '2025-01-03' +) +SELECT *, LAG(user, 1) OVER (PARTITION BY device ORDER BY "timestamp") AS prev +FROM events +WHERE device IS NOT NULL AND device != '' AND user IS NOT NULL AND user != '' +LIMIT 100; +---- +logical_plan +01)Projection: events.device, events.user, events.timestamp, lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS prev +02)--Limit: skip=0, fetch=100 +03)----WindowAggr: windowExpr=[[lag(events.user, Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] +04)------SubqueryAlias: events +05)--------Projection: __datafusion_extracted_1 AS device, __datafusion_extracted_2 AS user, events_mem.timestamp +06)----------Filter: events_mem.date = Date32("2025-01-03") AND __datafusion_extracted_1 IS NOT NULL AND __datafusion_extracted_1 != Utf8View("") AND __datafusion_extracted_2 IS NOT NULL AND __datafusion_extracted_2 != Utf8View("") +07)------------Projection: get_field(events_mem.ids, Utf8("id1")) AS __datafusion_extracted_1, get_field(events_mem.structs, Utf8("var1")) AS __datafusion_extracted_2, events_mem.date, events_mem.timestamp +08)--------------TableScan: events_mem projection=[date, timestamp, ids, structs] +physical_plan +01)ProjectionExec: expr=[device@0 as device, user@1 as user, timestamp@2 as timestamp, lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@3 as prev] +02)--GlobalLimitExec: skip=0, fetch=100 +03)----BoundedWindowAggExec: wdw=[lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": nullable Utf8View }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted] +04)------SortExec: expr=[device@0 ASC NULLS LAST, timestamp@2 ASC NULLS LAST], preserve_partitioning=[false] +05)--------ProjectionExec: expr=[__datafusion_extracted_1@0 as device, __datafusion_extracted_2@1 as user, timestamp@2 as timestamp] +06)----------FilterExec: date@2 = 2025-01-03 AND __datafusion_extracted_1@0 IS NOT NULL AND __datafusion_extracted_1@0 != AND __datafusion_extracted_2@1 IS NOT NULL AND __datafusion_extracted_2@1 != , projection=[__datafusion_extracted_1@0, __datafusion_extracted_2@1, timestamp@3] +07)------------ProjectionExec: expr=[get_field(ids@2, id1) as __datafusion_extracted_1, get_field(structs@3, var1) as __datafusion_extracted_2, date@0 as date, timestamp@1 as timestamp] +08)--------------DataSourceExec: partitions=1, partition_sizes=[1] + +# Determinism: the plan is a fixed point reached before the default pass limit, +# so it does not depend on which of the two rules runs last. One pass is enough +# for the simple shape. +statement ok +SET datafusion.optimizer.max_passes = 1; + +query TT +EXPLAIN SELECT ids['id1'] FROM events_mem WHERE "date" = '2025-01-03'; +---- +logical_plan +01)Projection: __datafusion_extracted_1 AS events_mem.ids[id1] +02)--Filter: events_mem.date = Date32("2025-01-03") +03)----Projection: get_field(events_mem.ids, Utf8("id1")) AS __datafusion_extracted_1, events_mem.date +04)------TableScan: events_mem projection=[date, ids] +physical_plan +01)ProjectionExec: expr=[__datafusion_extracted_1@0 as events_mem.ids[id1]] +02)--FilterExec: date@1 = 2025-01-03, projection=[__datafusion_extracted_1@0] +03)----ProjectionExec: expr=[get_field(ids@1, id1) as __datafusion_extracted_1, date@0 as date] +04)------DataSourceExec: partitions=1, partition_sizes=[1] + +statement ok +RESET datafusion.optimizer.max_passes; + +query TT +EXPLAIN SELECT ids['id1'] FROM events_mem WHERE "date" = '2025-01-03'; +---- +logical_plan +01)Projection: __datafusion_extracted_1 AS events_mem.ids[id1] +02)--Filter: events_mem.date = Date32("2025-01-03") +03)----Projection: get_field(events_mem.ids, Utf8("id1")) AS __datafusion_extracted_1, events_mem.date +04)------TableScan: events_mem projection=[date, ids] +physical_plan +01)ProjectionExec: expr=[__datafusion_extracted_1@0 as events_mem.ids[id1]] +02)--FilterExec: date@1 = 2025-01-03, projection=[__datafusion_extracted_1@0] +03)----ProjectionExec: expr=[get_field(ids@1, id1) as __datafusion_extracted_1, date@0 as date] +04)------DataSourceExec: partitions=1, partition_sizes=[1] + +# The #14540 shape needs a second pass, because the two filters merge into one. +# It is stable from that pass on. +statement ok +SET datafusion.optimizer.max_passes = 2; + +query TT +EXPLAIN WITH events AS ( + SELECT ids.id1 AS device, structs.var1 AS user, "timestamp" + FROM events_mem + WHERE "date" = '2025-01-03' +) +SELECT *, LAG(user, 1) OVER (PARTITION BY device ORDER BY "timestamp") AS prev +FROM events +WHERE device IS NOT NULL AND device != '' AND user IS NOT NULL AND user != '' +LIMIT 100; +---- +logical_plan +01)Projection: events.device, events.user, events.timestamp, lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS prev +02)--Limit: skip=0, fetch=100 +03)----WindowAggr: windowExpr=[[lag(events.user, Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] +04)------SubqueryAlias: events +05)--------Projection: __datafusion_extracted_1 AS device, __datafusion_extracted_2 AS user, events_mem.timestamp +06)----------Filter: events_mem.date = Date32("2025-01-03") AND __datafusion_extracted_1 IS NOT NULL AND __datafusion_extracted_1 != Utf8View("") AND __datafusion_extracted_2 IS NOT NULL AND __datafusion_extracted_2 != Utf8View("") +07)------------Projection: get_field(events_mem.ids, Utf8("id1")) AS __datafusion_extracted_1, get_field(events_mem.structs, Utf8("var1")) AS __datafusion_extracted_2, events_mem.date, events_mem.timestamp +08)--------------TableScan: events_mem projection=[date, timestamp, ids, structs] +physical_plan +01)ProjectionExec: expr=[device@0 as device, user@1 as user, timestamp@2 as timestamp, lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@3 as prev] +02)--GlobalLimitExec: skip=0, fetch=100 +03)----BoundedWindowAggExec: wdw=[lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": nullable Utf8View }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted] +04)------SortExec: expr=[device@0 ASC NULLS LAST, timestamp@2 ASC NULLS LAST], preserve_partitioning=[false] +05)--------ProjectionExec: expr=[__datafusion_extracted_1@0 as device, __datafusion_extracted_2@1 as user, timestamp@2 as timestamp] +06)----------FilterExec: date@2 = 2025-01-03 AND __datafusion_extracted_1@0 IS NOT NULL AND __datafusion_extracted_1@0 != AND __datafusion_extracted_2@1 IS NOT NULL AND __datafusion_extracted_2@1 != , projection=[__datafusion_extracted_1@0, __datafusion_extracted_2@1, timestamp@3] +07)------------ProjectionExec: expr=[get_field(ids@2, id1) as __datafusion_extracted_1, get_field(structs@3, var1) as __datafusion_extracted_2, date@0 as date, timestamp@1 as timestamp] +08)--------------DataSourceExec: partitions=1, partition_sizes=[1] + +statement ok +RESET datafusion.optimizer.max_passes; + +query TTT +WITH events AS ( + SELECT ids.id1 AS device, structs.var1 AS user, "timestamp" + FROM events_mem + WHERE "date" = '2025-01-03' +) +SELECT device, user, CAST(LAG(user, 1) OVER (PARTITION BY device ORDER BY "timestamp") AS VARCHAR) AS prev +FROM events +WHERE device IS NOT NULL AND device != '' AND user IS NOT NULL AND user != '' +ORDER BY user; +---- +dev1 user1 NULL +dev1 user2 user1 + +statement ok +DROP TABLE events_mem; + +# Parquet: both properties must survive. `DataSourceExec` reads only the struct +# leaf `ids.id1` (the `get_field` is in the scan projection), and the `date` +# predicate reaches the scan for row group pruning. +statement ok +COPY ( + SELECT + column1 AS "date", + column2 AS "timestamp", + column3 AS ids, + column4 AS structs + FROM VALUES + (DATE '2025-01-03', TIMESTAMP '2025-01-03 01:00:00', {id1: 'dev1', extra: 1}, {var1: 'user1', extra: 'e1'}), + (DATE '2025-01-03', TIMESTAMP '2025-01-03 02:00:00', {id1: 'dev1', extra: 2}, {var1: 'user2', extra: 'e2'}), + (DATE '2025-01-04', TIMESTAMP '2025-01-04 01:00:00', {id1: 'dev2', extra: 3}, {var1: 'user3', extra: 'e3'}) +) TO 'test_files/scratch/projection_pushdown/events.parquet' +STORED AS PARQUET; + +statement ok +CREATE EXTERNAL TABLE events_pq STORED AS PARQUET +LOCATION 'test_files/scratch/projection_pushdown/events.parquet'; + +query TT +EXPLAIN SELECT ids['id1'] FROM events_pq WHERE "date" = '2025-01-03'; +---- +logical_plan +01)Projection: __datafusion_extracted_1 AS events_pq.ids[id1] +02)--Filter: events_pq.date = Date32("2025-01-03") +03)----Projection: get_field(events_pq.ids, Utf8("id1")) AS __datafusion_extracted_1, events_pq.date +04)------TableScan: events_pq projection=[date, ids], partial_filters=[events_pq.date = Date32("2025-01-03")] +physical_plan +01)ProjectionExec: expr=[__datafusion_extracted_1@0 as events_pq.ids[id1]] +02)--FilterExec: date@1 = 2025-01-03, projection=[__datafusion_extracted_1@0] +03)----DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/projection_pushdown/events.parquet]]}, projection=[get_field(ids@2, id1) as __datafusion_extracted_1, date], file_type=parquet, predicate=date@0 = 2025-01-03, pruning_predicate=date_null_count@2 != row_count@3 AND date_min@0 <= 2025-01-03 AND 2025-01-03 <= date_max@1, required_guarantees=[date in (2025-01-03)] + +query TT +EXPLAIN WITH events AS ( + SELECT ids.id1 AS device, structs.var1 AS user, "timestamp" + FROM events_pq + WHERE "date" = '2025-01-03' +) +SELECT *, LAG(user, 1) OVER (PARTITION BY device ORDER BY "timestamp") AS prev +FROM events +WHERE device IS NOT NULL AND device != '' AND user IS NOT NULL AND user != '' +LIMIT 100; +---- +logical_plan +01)Projection: events.device, events.user, events.timestamp, lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS prev +02)--Limit: skip=0, fetch=100 +03)----WindowAggr: windowExpr=[[lag(events.user, Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] +04)------SubqueryAlias: events +05)--------Projection: __datafusion_extracted_1 AS device, __datafusion_extracted_2 AS user, events_pq.timestamp +06)----------Filter: events_pq.date = Date32("2025-01-03") AND __datafusion_extracted_1 IS NOT NULL AND __datafusion_extracted_1 != Utf8("") AND __datafusion_extracted_2 IS NOT NULL AND __datafusion_extracted_2 != Utf8("") +07)------------Projection: get_field(events_pq.ids, Utf8("id1")) AS __datafusion_extracted_1, get_field(events_pq.structs, Utf8("var1")) AS __datafusion_extracted_2, events_pq.date, events_pq.timestamp +08)--------------TableScan: events_pq projection=[date, timestamp, ids, structs], partial_filters=[events_pq.date = Date32("2025-01-03")] +physical_plan +01)ProjectionExec: expr=[device@0 as device, user@1 as user, timestamp@2 as timestamp, lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@3 as prev] +02)--GlobalLimitExec: skip=0, fetch=100 +03)----BoundedWindowAggExec: wdw=[lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": nullable Utf8 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted] +04)------SortExec: expr=[device@0 ASC NULLS LAST, timestamp@2 ASC NULLS LAST], preserve_partitioning=[false] +05)--------ProjectionExec: expr=[__datafusion_extracted_1@0 as device, __datafusion_extracted_2@1 as user, timestamp@2 as timestamp] +06)----------FilterExec: date@2 = 2025-01-03 AND __datafusion_extracted_1@0 IS NOT NULL AND __datafusion_extracted_1@0 != AND __datafusion_extracted_2@1 IS NOT NULL AND __datafusion_extracted_2@1 != , projection=[__datafusion_extracted_1@0, __datafusion_extracted_2@1, timestamp@3] +07)------------DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/projection_pushdown/events.parquet]]}, projection=[get_field(ids@2, id1) as __datafusion_extracted_1, get_field(structs@3, var1) as __datafusion_extracted_2, date, timestamp], file_type=parquet, predicate=date@0 = 2025-01-03 AND get_field(ids@2, id1) IS NOT NULL AND get_field(ids@2, id1) != AND get_field(structs@3, var1) IS NOT NULL AND get_field(structs@3, var1) != , pruning_predicate=date_null_count@2 != row_count@3 AND date_min@0 <= 2025-01-03 AND 2025-01-03 <= date_max@1, required_guarantees=[date in (2025-01-03)] + +statement ok +DROP TABLE events_pq; + +statement ok +RESET datafusion.optimizer.max_passes; + +statement ok +SET datafusion.execution.target_partitions = 4; From 406e428657115bbd268bf1ed13bd746a3080835e Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Tue, 6 Oct 2026 09:19:29 -0400 Subject: [PATCH 4/4] test: check per rule that the pushdown rules reach a fixed point Address review: add an observer test that fails when PushDownFilter or PushDownLeafProjections changes the plan in the last optimizer pass, and correct the rule precedence wording in the optimizer guide. Co-Authored-By: Claude Opus 5.5 --- datafusion/optimizer/src/push_down_filter.rs | 59 +++++++++++++++++++ .../library-user-guide/query-optimizer.md | 4 +- 2 files changed, 61 insertions(+), 2 deletions(-) diff --git a/datafusion/optimizer/src/push_down_filter.rs b/datafusion/optimizer/src/push_down_filter.rs index 370466f8ee971..33664935233bf 100644 --- a/datafusion/optimizer/src/push_down_filter.rs +++ b/datafusion/optimizer/src/push_down_filter.rs @@ -4627,6 +4627,65 @@ mod tests { ) } + /// Runs the default optimizer and returns, for each pass, the names of the + /// rules that changed the plan in that pass. + fn rules_that_changed_plan_per_pass(plan: LogicalPlan) -> Result>> { + let optimizer = Optimizer::new(); + let rules_per_pass = optimizer.rules.len(); + let mut previous = plan.display_indent().to_string(); + let mut calls = 0; + let mut passes: Vec> = vec![]; + optimizer.optimize(plan, &OptimizerContext::new(), |plan, rule| { + if calls % rules_per_pass == 0 { + passes.push(vec![]); + } + calls += 1; + let current = plan.display_indent().to_string(); + if current != previous { + passes.last_mut().unwrap().push(rule.name().to_string()); + previous = current; + } + })?; + Ok(passes) + } + + /// `PushDownFilter` and `PushDownLeafProjections` must not undo each other + /// for a filter next to a pure extraction projection + /// (). Neither rule may + /// change the plan in the last optimizer pass. The source does not absorb + /// filters, so the `Filter` node stays in the plan. + #[test] + fn filter_and_extraction_projection_reach_fixed_point() -> Result<()> { + let scan = || { + table_scan_with_pushdown_provider_builder( + TableProviderFilterPushDown::Unsupported, + vec![], + None, + ) + }; + let simple = scan()? + .filter(col("b").gt(lit(5)))? + .project(vec![leaf_udf_expr(col("a"))])? + .build()?; + // Two filters, one of them on an extracted leaf. + let two_filters = scan()? + .filter(col("b").gt(lit(5)))? + .filter(leaf_udf_expr(col("a")).eq(lit(1)))? + .project(vec![leaf_udf_expr(col("a")), col("b")])? + .build()?; + + for plan in [simple, two_filters] { + let passes = rules_that_changed_plan_per_pass(plan)?; + let last = passes.last().unwrap(); + assert!( + !last.iter().any(|rule| rule == "push_down_filter" + || rule == "push_down_leaf_projections"), + "pushdown rules changed the plan in the last pass: {passes:?}" + ); + } + Ok(()) + } + /// A projection that mixes an extraction alias with a computed expression is /// not a pure extraction projection, so the filter still moves below it. #[test] diff --git a/docs/source/library-user-guide/query-optimizer.md b/docs/source/library-user-guide/query-optimizer.md index 883f922c06df9..73c3ed6ed11dc 100644 --- a/docs/source/library-user-guide/query-optimizer.md +++ b/docs/source/library-user-guide/query-optimizer.md @@ -164,9 +164,9 @@ the actual optimization. This approach is used in projection push down and filte ### Rule Precedence -Two rules can want the opposite order for the same pair of adjacent plan nodes. Each rule then undoes the work of the other one on every optimizer pass. The loop stops at `datafusion.optimizer.max_passes`, and the position of the two rules in the rule list decides the plan. This wastes a pass and makes the plan depend on the rule list. +Two rules can want the opposite order for the same pair of adjacent plan nodes. Each rule then undoes the work of the other one on every optimizer pass. The plan does not reach a fixed point, so the optimizer runs more passes than necessary, and the position of the two rules in the rule list decides the plan. The optimizer can stop early when the plan at the end of a pass is the same as the plan at the end of an earlier pass, but that does not make the result independent of the rule order. -Do not let two rules compete. Give one rule precedence, make that rule yield, and record the decision in the module documentation of both rules. +Do not let two rules compete. Give one rule precedence, make the competing rule yield, and record the decision in the module documentation of both rules. There is one such decision today: