From c6be450f4fe0adb357d3487b75e2dfb4bbdb0779 Mon Sep 17 00:00:00 2001 From: Qi Zhu Date: Tue, 29 Sep 2026 14:34:27 +0800 Subject: [PATCH 1/6] fix: keep a volatile predicate above a window function The window arm of push_down_filter pushed any predicate whose column refs were all partition keys. A volatile predicate that reads no columns at all, such as `random() < 0.5`, satisfies that test vacuously and was pushed below the window, where it changes which rows the window function sees and so the value it computes for the rows that survive. The aggregate arm already drops volatile group expressions before the same kind of check; do the equivalent here and keep volatile predicates above the window. --- datafusion/optimizer/src/push_down_filter.rs | 47 +++++++++++++++++++- 1 file changed, 46 insertions(+), 1 deletion(-) diff --git a/datafusion/optimizer/src/push_down_filter.rs b/datafusion/optimizer/src/push_down_filter.rs index 1972174a45f8c..d70130bbadd17 100644 --- a/datafusion/optimizer/src/push_down_filter.rs +++ b/datafusion/optimizer/src/push_down_filter.rs @@ -1109,7 +1109,15 @@ impl OptimizerRule for PushDownFilter { let mut push_predicates = vec![]; for expr in predicates { let cols = expr.column_refs(); - if cols.iter().all(|c| potential_partition_keys.contains(c)) { + // A volatile predicate has to stay above the window. Pushing it + // changes which rows the window function sees, and so the value + // it computes for the rows that do survive. Checking this first + // also covers a volatile predicate that reads no columns at all, + // such as `random() < 0.5`, which would otherwise satisfy the + // partition-key test vacuously. + if !expr.is_volatile() + && cols.iter().all(|c| potential_partition_keys.contains(c)) + { push_predicates.push(expr); } else { keep_predicates.push(expr); @@ -1924,6 +1932,43 @@ mod tests { ) } + /// verifies that a volatile predicate is not pushed through a window, even + /// when it reads no columns and so trivially satisfies the partition-key test + #[test] + fn filter_volatile_keep_window() -> Result<()> { + let table_scan = test_table_scan()?; + + let window = Expr::from(WindowFunction::new( + WindowFunctionDefinition::WindowUDF( + datafusion_functions_window::rank::rank_udwf(), + ), + vec![], + )) + .partition_by(vec![col("a")]) + .order_by(vec![col("c").sort(true, true)]) + .build() + .unwrap(); + + let fun = ScalarUDF::new_from_impl(TestScalarUDF { + signature: Signature::exact(vec![], Volatility::Volatile), + }); + let volatile = Expr::ScalarFunction(ScalarFunction::new_udf(Arc::new(fun), vec![])); + + let plan = LogicalPlanBuilder::from(table_scan) + .window(vec![window])? + .filter(volatile.gt(lit(10i64)))? + .build()?; + + assert_optimized_plan_equal!( + plan, + @r" + Filter: TestScalarUDF() > Int64(10) + WindowAggr: windowExpr=[[rank() PARTITION BY [test.a] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] + TableScan: test + " + ) + } + /// verifies that filters are not pushed on order by columns (that are not used in partitioning) #[test] fn filter_order_keep_window() -> Result<()> { From 4dc3f4c9a77bf40935eee633a8c1ed6a69a48337 Mon Sep 17 00:00:00 2001 From: Qi Zhu Date: Tue, 29 Sep 2026 14:39:31 +0800 Subject: [PATCH 2/6] feat: push window predicates on expression PARTITION BY keys A predicate that depends only on a window's PARTITION BY keys is constant within every partition, so applying it below the window drops whole partitions and leaves every surviving row's window value unchanged. That never fired for expression keys such as `PARTITION BY a + b`, because each key was mapped through `qualified_name()` into a column literally named "a + b", which the predicate's real column refs could never match. Match each conjunct against the key expressions themselves instead: walk the predicate, treat a subtree equal to one of the keys as satisfied in full, and reject only on reaching a Column no key covered. This is a strict generalization, since for a plain column key the reference to that column is itself a subtree equal to the key. It also fixes mixed predicates. With PARTITION BY year, num * num the predicate `year = '2021' OR num * num > 4` is constant within every partition, but its column refs are {year, num} against a key set of {year, "num * num"}, so num was not found and the whole conjunct stayed above the window. Matching is structural, so this stays conservative: a predicate written `b + a` does not match a key written `a + b` and is left in place. --- datafusion/optimizer/src/push_down_filter.rs | 136 ++++++++++++++++--- 1 file changed, 119 insertions(+), 17 deletions(-) diff --git a/datafusion/optimizer/src/push_down_filter.rs b/datafusion/optimizer/src/push_down_filter.rs index d70130bbadd17..bfc6caf60fd9d 100644 --- a/datafusion/optimizer/src/push_down_filter.rs +++ b/datafusion/optimizer/src/push_down_filter.rs @@ -26,7 +26,9 @@ use itertools::Itertools; use log::{Level, debug, log_enabled}; use datafusion_common::instant::Instant; -use datafusion_common::tree_node::{Transformed, TransformedResult, TreeNode}; +use datafusion_common::tree_node::{ + Transformed, TransformedResult, TreeNode, TreeNodeRecursion, +}; use datafusion_common::{ Column, DFSchema, Result, assert_eq_or_internal_err, internal_err, plan_err, qualified_name, @@ -1071,8 +1073,13 @@ impl OptimizerRule for PushDownFilter { // multiple window functions, each with potentially different partition keys. // Therefore, we need to ensure that any potential partition key returned is used in // ALL window functions. Otherwise, filters cannot be pushed by through that column. - fn extract_partition_keys(func: &WindowFunction) -> HashSet { - expr_columns(&func.params.partition_by) + // Keyed by the partition *expression*, not by a name synthesised + // from it. `PARTITION BY a + b` used to be mapped through + // `qualified_name()` to a column literally called "a + b", which no + // predicate's real column refs could ever match, so such a key was + // dead weight in this set. + fn extract_partition_keys(func: &WindowFunction) -> HashSet { + func.params.partition_by.iter().cloned().collect() } let potential_partition_keys = window @@ -1108,7 +1115,6 @@ impl OptimizerRule for PushDownFilter { let mut keep_predicates = vec![]; let mut push_predicates = vec![]; for expr in predicates { - let cols = expr.column_refs(); // A volatile predicate has to stay above the window. Pushing it // changes which rows the window function sees, and so the value // it computes for the rows that do survive. Checking this first @@ -1116,7 +1122,7 @@ impl OptimizerRule for PushDownFilter { // such as `random() < 0.5`, which would otherwise satisfy the // partition-key test vacuously. if !expr.is_volatile() - && cols.iter().all(|c| potential_partition_keys.contains(c)) + && reads_only_partition_keys(&expr, &potential_partition_keys)? { push_predicates.push(expr); } else { @@ -1125,12 +1131,11 @@ impl OptimizerRule for PushDownFilter { } // Unlike with aggregations, there are no cases where we have to replace, e.g., - // `a+b` with Column(a)+Column(b). This is because partition expressions are not - // available as standalone columns to the user. For example, while an aggregation on - // `a+b` becomes Column(a + b), in a window partition it becomes - // `func() PARTITION BY [a + b] ...`. Thus, filters on expressions always remain in - // place, so we can use `push_predicates` directly. This is consistent with other - // optimizers, such as the one used by Postgres. + // `a+b` with Column(a+b). This is because partition expressions are not available + // as standalone columns to the user: while an aggregation on `a+b` becomes + // Column(a + b), in a window partition it stays `func() PARTITION BY [a + b] ...`. + // That is why the predicate is matched against the key expressions themselves and + // can be pushed unchanged. // If we have a filter to push, we push it down to the input of the aggregate let result = if let Some(predicate) = conjunction(push_predicates) { @@ -1527,6 +1532,37 @@ fn with_filters(predicates: Vec, plan: LogicalPlan) -> LogicalPlan { } } +/// Does `expr` read nothing beyond the given window partition keys? +/// +/// A subtree that is exactly one of the keys counts as read in full, so a +/// predicate on an *expression* key, say `NULLIF(c, '') IS NOT NULL` against +/// `PARTITION BY NULLIF(c, '')`, qualifies even though the column it ultimately +/// reads (`c`) is not a key on its own. Such a predicate is constant within each +/// partition, so applying it below the window drops whole partitions and leaves +/// every surviving row's window value unchanged. +/// +/// Matching is structural, which makes this conservative rather than wrong: a +/// predicate written `b + a` does not match a key written `a + b`, and is simply +/// left above the window. +fn reads_only_partition_keys( + expr: &Expr, + partition_keys: &HashSet, +) -> Result { + let mut reads_something_else = false; + expr.apply(|node| { + Ok(if partition_keys.contains(node) { + // the whole key was matched, so whatever it reads is accounted for + TreeNodeRecursion::Jump + } else if matches!(node, Expr::Column(_)) { + reads_something_else = true; + TreeNodeRecursion::Stop + } else { + TreeNodeRecursion::Continue + }) + })?; + Ok(!reads_something_else) +} + fn expr_columns(exprs: &[Expr]) -> HashSet { exprs .iter() @@ -1898,10 +1934,10 @@ mod tests { ) } - /// verifies that filters on partition expressions are not pushed, as the single expression - /// column is not available to the user, unlike with aggregations + /// verifies that a filter on an expression partition key is pushed, matched + /// against the key expression itself rather than against its column refs #[test] - fn filter_expression_keep_window() -> Result<()> { + fn filter_expression_move_window() -> Result<()> { let table_scan = test_table_scan()?; let window = Expr::from(WindowFunction::new( @@ -1917,15 +1953,81 @@ mod tests { let plan = LogicalPlanBuilder::from(table_scan) .window(vec![window])? - // unlike with aggregations, single partition column "test.a + test.b" is not available - // to the plan, so we use multiple columns when filtering + // the single partition column "test.a + test.b" is not available to the + // plan, so the predicate is written over the underlying columns .filter(add(col("a"), col("b")).gt(lit(10i64)))? .build()?; assert_optimized_plan_equal!( plan, @r" - Filter: test.a + test.b > Int64(10) + WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] + TableScan: test, full_filters=[test.a + test.b > Int64(10)] + " + ) + } + + /// verifies that a single predicate spanning a column key and an expression + /// key is pushed: it is constant within every partition, but neither the old + /// column-name matching nor a rule that looked only at expression keys could + /// see that + #[test] + fn filter_move_window_mixed_column_and_expression_keys() -> Result<()> { + let table_scan = test_table_scan()?; + + let window = Expr::from(WindowFunction::new( + WindowFunctionDefinition::WindowUDF( + datafusion_functions_window::rank::rank_udwf(), + ), + vec![], + )) + .partition_by(vec![col("a"), add(col("a"), col("b"))]) // PARTITION BY a, a + b + .order_by(vec![col("c").sort(true, true)]) + .build() + .unwrap(); + + let plan = LogicalPlanBuilder::from(table_scan) + .window(vec![window])? + .filter(col("a").gt(lit(1i64)).or(add(col("a"), col("b")).gt(lit(10i64))))? + .build()?; + + assert_optimized_plan_equal!( + plan, + @r" + WindowAggr: windowExpr=[[rank() PARTITION BY [test.a, test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] + TableScan: test, full_filters=[test.a > Int64(1) OR test.a + test.b > Int64(10)] + " + ) + } + + /// verifies that an expression partition key does not make the columns it + /// reads pushable on their own: `a` alone is not constant within a partition + /// of `a + b` + #[test] + fn filter_keep_window_column_underlying_expression_key() -> Result<()> { + let table_scan = test_table_scan()?; + + let window = Expr::from(WindowFunction::new( + WindowFunctionDefinition::WindowUDF( + datafusion_functions_window::rank::rank_udwf(), + ), + vec![], + )) + .partition_by(vec![add(col("a"), col("b"))]) // PARTITION BY a + b + .order_by(vec![col("c").sort(true, true)]) + .build() + .unwrap(); + + let plan = LogicalPlanBuilder::from(table_scan) + .window(vec![window])? + .filter(col("a").gt(lit(10i64)))? + .build()?; + assert_plan_not_transformed!(plan.clone()); + + assert_optimized_plan_equal!( + plan, + @r" + Filter: test.a > Int64(10) WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] TableScan: test " From 5acfd50c871dbd53ad3a3571ba4d077deb747484 Mon Sep 17 00:00:00 2001 From: Qi Zhu Date: Wed, 30 Sep 2026 10:40:38 +0800 Subject: [PATCH 3/6] test: cover expression partition keys in window filter pushdown Also run rustfmt over the file, which reflows one filter built in an earlier commit of this branch. The new cases are the ones the change turns on or deliberately leaves alone, and none of them were reachable from the existing tests: - a function expression key, the shape the rule exists for, where the key is not arithmetic and the column it reads is not a key by itself - a predicate written with the operands in the other order, which does not match structurally and so stays above the window - a predicate built on top of a key, where the key subtree accounts for everything read and the surrounding literal reads nothing - an expression key shared by every window function, and one present in only one of them - a volatile predicate whose operand matches a key, to pin that the volatile test still runs first --- datafusion/optimizer/src/push_down_filter.rs | 243 ++++++++++++++++++- 1 file changed, 241 insertions(+), 2 deletions(-) diff --git a/datafusion/optimizer/src/push_down_filter.rs b/datafusion/optimizer/src/push_down_filter.rs index bfc6caf60fd9d..7e29afd61c74b 100644 --- a/datafusion/optimizer/src/push_down_filter.rs +++ b/datafusion/optimizer/src/push_down_filter.rs @@ -1988,7 +1988,11 @@ mod tests { let plan = LogicalPlanBuilder::from(table_scan) .window(vec![window])? - .filter(col("a").gt(lit(1i64)).or(add(col("a"), col("b")).gt(lit(10i64))))? + .filter( + col("a") + .gt(lit(1i64)) + .or(add(col("a"), col("b")).gt(lit(10i64))), + )? .build()?; assert_optimized_plan_equal!( @@ -2034,6 +2038,240 @@ mod tests { ) } + /// verifies that a predicate on a *function* expression key is pushed. This is + /// the shape the rule exists for: `PARTITION BY NULLIF(c, '')` with a predicate + /// on `NULLIF(c, '')`, where the key is not an arithmetic expression and the + /// column it reads is not a key on its own. + #[test] + fn filter_move_window_function_expression_key() -> Result<()> { + let table_scan = test_table_scan()?; + + let window = Expr::from(WindowFunction::new( + WindowFunctionDefinition::WindowUDF( + datafusion_functions_window::rank::rank_udwf(), + ), + vec![], + )) + .partition_by(vec![immutable_udf_call(col("c"))]) + .order_by(vec![col("c").sort(true, true)]) + .build() + .unwrap(); + + let plan = LogicalPlanBuilder::from(table_scan) + .window(vec![window])? + .filter(immutable_udf_call(col("c")).is_not_null())? + .build()?; + + assert_optimized_plan_equal!( + plan, + @r" + WindowAggr: windowExpr=[[rank() PARTITION BY [TestScalarUDF(test.c)] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] + TableScan: test, full_filters=[TestScalarUDF(test.c) IS NOT NULL] + " + ) + } + + /// verifies that matching is structural, so a predicate written with the + /// operands in the other order does not match the key and stays above the + /// window. Conservative rather than wrong. + #[test] + fn filter_keep_window_expression_key_operand_order() -> Result<()> { + let table_scan = test_table_scan()?; + + let window = Expr::from(WindowFunction::new( + WindowFunctionDefinition::WindowUDF( + datafusion_functions_window::rank::rank_udwf(), + ), + vec![], + )) + .partition_by(vec![add(col("a"), col("b"))]) // PARTITION BY a + b + .order_by(vec![col("c").sort(true, true)]) + .build() + .unwrap(); + + let plan = LogicalPlanBuilder::from(table_scan) + .window(vec![window])? + .filter(add(col("b"), col("a")).gt(lit(10i64)))? // b + a, not a + b + .build()?; + assert_plan_not_transformed!(plan.clone()); + + assert_optimized_plan_equal!( + plan, + @r" + Filter: test.b + test.a > Int64(10) + WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] + TableScan: test + " + ) + } + + /// verifies that a predicate built *on top of* an expression key is pushed: + /// the key subtree accounts for everything it reads, and the literal around it + /// reads nothing, so the predicate is still constant within each partition + #[test] + fn filter_move_window_expression_key_subtree() -> Result<()> { + let table_scan = test_table_scan()?; + + let window = Expr::from(WindowFunction::new( + WindowFunctionDefinition::WindowUDF( + datafusion_functions_window::rank::rank_udwf(), + ), + vec![], + )) + .partition_by(vec![add(col("a"), col("b"))]) // PARTITION BY a + b + .order_by(vec![col("c").sort(true, true)]) + .build() + .unwrap(); + + let plan = LogicalPlanBuilder::from(table_scan) + .window(vec![window])? + .filter(add(add(col("a"), col("b")), lit(1i64)).gt(lit(10i64)))? + .build()?; + + assert_optimized_plan_equal!( + plan, + @r" + WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] + TableScan: test, full_filters=[test.a + test.b + Int64(1) > Int64(10)] + " + ) + } + + /// verifies that an expression key shared by every window function is pushed, + /// the expression-key counterpart of filter_multiple_windows_common_partitions + #[test] + fn filter_multiple_windows_common_expression_partition() -> Result<()> { + let table_scan = test_table_scan()?; + + let window1 = Expr::from(WindowFunction::new( + WindowFunctionDefinition::WindowUDF( + datafusion_functions_window::rank::rank_udwf(), + ), + vec![], + )) + .partition_by(vec![add(col("a"), col("b"))]) + .order_by(vec![col("c").sort(true, true)]) + .build() + .unwrap(); + + let window2 = Expr::from(WindowFunction::new( + WindowFunctionDefinition::WindowUDF( + datafusion_functions_window::rank::rank_udwf(), + ), + vec![], + )) + .partition_by(vec![col("b"), add(col("a"), col("b"))]) + .order_by(vec![col("c").sort(true, true)]) + .build() + .unwrap(); + + let plan = LogicalPlanBuilder::from(table_scan) + .window(vec![window1, window2])? + .filter(add(col("a"), col("b")).gt(lit(10i64)))? // a + b is in both + .build()?; + + assert_optimized_plan_equal!( + plan, + @r" + WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, rank() PARTITION BY [test.b, test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] + TableScan: test, full_filters=[test.a + test.b > Int64(10)] + " + ) + } + + /// verifies that an expression key present in only one of the window functions + /// is not pushed, since the other function's partitions would change + #[test] + fn filter_multiple_windows_disjoint_expression_partition() -> Result<()> { + let table_scan = test_table_scan()?; + + let window1 = Expr::from(WindowFunction::new( + WindowFunctionDefinition::WindowUDF( + datafusion_functions_window::rank::rank_udwf(), + ), + vec![], + )) + .partition_by(vec![add(col("a"), col("b"))]) + .order_by(vec![col("c").sort(true, true)]) + .build() + .unwrap(); + + let window2 = Expr::from(WindowFunction::new( + WindowFunctionDefinition::WindowUDF( + datafusion_functions_window::rank::rank_udwf(), + ), + vec![], + )) + .partition_by(vec![col("a")]) + .order_by(vec![col("c").sort(true, true)]) + .build() + .unwrap(); + + let plan = LogicalPlanBuilder::from(table_scan) + .window(vec![window1, window2])? + .filter(add(col("a"), col("b")).gt(lit(10i64)))? // a + b is in one only + .build()?; + assert_plan_not_transformed!(plan.clone()); + + assert_optimized_plan_equal!( + plan, + @r" + Filter: test.a + test.b > Int64(10) + WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, rank() PARTITION BY [test.a] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] + TableScan: test + " + ) + } + + /// verifies that the volatile check still wins over an expression key: the + /// predicate matches the key structurally, but pushing it would change which + /// rows each partition contains and so the window values of the survivors + #[test] + fn filter_volatile_keep_window_expression_key() -> Result<()> { + let table_scan = test_table_scan()?; + + let window = Expr::from(WindowFunction::new( + WindowFunctionDefinition::WindowUDF( + datafusion_functions_window::rank::rank_udwf(), + ), + vec![], + )) + .partition_by(vec![add(col("a"), col("b"))]) + .order_by(vec![col("c").sort(true, true)]) + .build() + .unwrap(); + + let fun = ScalarUDF::new_from_impl(TestScalarUDF { + signature: Signature::exact(vec![], Volatility::Volatile), + }); + let volatile = + Expr::ScalarFunction(ScalarFunction::new_udf(Arc::new(fun), vec![])); + + let plan = LogicalPlanBuilder::from(table_scan) + .window(vec![window])? + .filter(add(col("a"), col("b")).gt(volatile))? + .build()?; + assert_plan_not_transformed!(plan.clone()); + + assert_optimized_plan_equal!( + plan, + @r" + Filter: test.a + test.b > TestScalarUDF() + WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] + TableScan: test + " + ) + } + + /// An immutable one-argument UDF call, for building an expression partition + /// key that is a function call rather than an arithmetic expression. + fn immutable_udf_call(arg: Expr) -> Expr { + let fun = ScalarUDF::new_from_impl(TestScalarUDF { + signature: Signature::any(1, Volatility::Immutable), + }); + Expr::ScalarFunction(ScalarFunction::new_udf(Arc::new(fun), vec![arg])) + } + /// verifies that a volatile predicate is not pushed through a window, even /// when it reads no columns and so trivially satisfies the partition-key test #[test] @@ -2054,7 +2292,8 @@ mod tests { let fun = ScalarUDF::new_from_impl(TestScalarUDF { signature: Signature::exact(vec![], Volatility::Volatile), }); - let volatile = Expr::ScalarFunction(ScalarFunction::new_udf(Arc::new(fun), vec![])); + let volatile = + Expr::ScalarFunction(ScalarFunction::new_udf(Arc::new(fun), vec![])); let plan = LogicalPlanBuilder::from(table_scan) .window(vec![window])? From eb779972d8b0fe4108b92454a04ff51a5d141bb7 Mon Sep 17 00:00:00 2001 From: Qi Zhu Date: Wed, 30 Sep 2026 11:21:28 +0800 Subject: [PATCH 4/6] fix: keep a subquery predicate above the window Reported by Copilot on the PR. `Expr::apply` does not enter a subquery's plan, so the columns it correlates on are invisible to the key test: `Exists` and `ScalarSubquery` are leaves, and `InSubquery` and `SetComparison` expose only their left-hand expression. With `PARTITION BY a + b` the predicate a + b > (SELECT ... WHERE inner.x = outer.a) therefore looked like it read nothing but the key, and was pushed. The outer reference varies inside a single partition, so the predicate is not constant there and pushing it changes the result. The same blind spot hides a volatile function inside a subquery from the `is_volatile` check that runs first. The previous column-key matching kept this expression above the window, since the predicate's column refs `{a, b}` never matched the synthesised key name. For a plain column key the hole predates this branch, but the expression-key matching would have widened it, so treat any subquery-bearing node as reading something the walk cannot account for. Two regression tests, one for the scalar subquery and one for `IN (subquery)`. Both fail without the guard, at the `assert_plan_not_transformed!` line. --- datafusion/optimizer/src/push_down_filter.rs | 116 ++++++++++++++++++- 1 file changed, 114 insertions(+), 2 deletions(-) diff --git a/datafusion/optimizer/src/push_down_filter.rs b/datafusion/optimizer/src/push_down_filter.rs index 7e29afd61c74b..f60ce32551d72 100644 --- a/datafusion/optimizer/src/push_down_filter.rs +++ b/datafusion/optimizer/src/push_down_filter.rs @@ -1544,6 +1544,14 @@ fn with_filters(predicates: Vec, plan: LogicalPlan) -> LogicalPlan { /// Matching is structural, which makes this conservative rather than wrong: a /// predicate written `b + a` does not match a key written `a + b`, and is simply /// left above the window. +/// +/// A node carrying a subquery counts as reading something else, whatever the +/// keys are. `Expr::apply` does not descend into a subquery's plan, so the +/// columns it correlates on are invisible here, and an outer reference can vary +/// inside a single partition: with `PARTITION BY a + b`, the predicate +/// `a + b > (SELECT ... WHERE inner.x = outer.a)` would otherwise look like it +/// reads nothing but the key. The same blind spot hides a volatile function +/// inside a subquery from the caller's `is_volatile` check. fn reads_only_partition_keys( expr: &Expr, partition_keys: &HashSet, @@ -1553,7 +1561,7 @@ fn reads_only_partition_keys( Ok(if partition_keys.contains(node) { // the whole key was matched, so whatever it reads is accounted for TreeNodeRecursion::Jump - } else if matches!(node, Expr::Column(_)) { + } else if reads_beyond_this_node(node) { reads_something_else = true; TreeNodeRecursion::Stop } else { @@ -1563,6 +1571,24 @@ fn reads_only_partition_keys( Ok(!reads_something_else) } +/// Does this node read data that walking its children cannot account for? +/// +/// A column reads itself. The subquery-bearing variants read whatever their +/// plan reads, including outer references to the window's own input, and +/// `Expr::apply` yields none of that: `Exists` and `ScalarSubquery` are leaves, +/// and `InSubquery` and `SetComparison` expose only their left-hand expression. +fn reads_beyond_this_node(node: &Expr) -> bool { + matches!( + node, + Expr::Column(_) + | Expr::OuterReferenceColumn(..) + | Expr::ScalarSubquery(_) + | Expr::Exists(_) + | Expr::InSubquery(_) + | Expr::SetComparison(_) + ) +} + fn expr_columns(exprs: &[Expr]) -> HashSet { exprs .iter() @@ -1588,7 +1614,7 @@ mod tests { ColumnarValue, ExprFunctionExt, Extension, LogicalPlanBuilder, ScalarFunctionArgs, ScalarUDF, ScalarUDFImpl, Signature, TableScan, TableSource, TableType, UserDefinedLogicalNodeCore, Volatility, WindowFunctionDefinition, col, - in_list, in_subquery, lit, + in_list, in_subquery, lit, out_ref_col, scalar_subquery, }; use crate::OptimizerContext; @@ -2263,6 +2289,92 @@ mod tests { ) } + /// verifies that a predicate carrying a correlated scalar subquery stays + /// above the window even though its visible part is exactly the expression + /// key. `Expr::apply` does not enter the subquery plan, so `outer.a`, which + /// varies inside a partition of `a + b`, is invisible to the key test. + #[test] + fn filter_keep_window_scalar_subquery_over_expression_key() -> Result<()> { + let table_scan = test_table_scan()?; + let subplan = Arc::new( + LogicalPlanBuilder::from(test_table_scan_with_name("sq")?) + .filter(col("sq.c").eq(out_ref_col(DataType::UInt32, "test.a")))? + .aggregate(Vec::::new(), vec![sum(col("sq.c"))])? + .build()?, + ); + + let window = Expr::from(WindowFunction::new( + WindowFunctionDefinition::WindowUDF( + datafusion_functions_window::rank::rank_udwf(), + ), + vec![], + )) + .partition_by(vec![add(col("a"), col("b"))]) // PARTITION BY a + b + .order_by(vec![col("c").sort(true, true)]) + .build() + .unwrap(); + + let plan = LogicalPlanBuilder::from(table_scan) + .window(vec![window])? + .filter(add(col("a"), col("b")).gt(scalar_subquery(subplan)))? + .build()?; + assert_plan_not_transformed!(plan.clone()); + + assert_optimized_plan_equal!( + plan, + @r" + Filter: test.a + test.b > () + Subquery: + Aggregate: groupBy=[[]], aggr=[[sum(sq.c)]] + TableScan: sq, full_filters=[sq.c = outer_ref(test.a)] + WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] + TableScan: test + " + ) + } + + /// the same guard for `IN (subquery)`, whose subquery plan `Expr::apply` + /// also skips: only the left-hand expression is walked + #[test] + fn filter_keep_window_in_subquery_over_expression_key() -> Result<()> { + let table_scan = test_table_scan()?; + let subplan = Arc::new( + LogicalPlanBuilder::from(test_table_scan_with_name("sq")?) + .filter(col("sq.c").eq(out_ref_col(DataType::UInt32, "test.a")))? + .project(vec![col("sq.c")])? + .build()?, + ); + + let window = Expr::from(WindowFunction::new( + WindowFunctionDefinition::WindowUDF( + datafusion_functions_window::rank::rank_udwf(), + ), + vec![], + )) + .partition_by(vec![add(col("a"), col("b"))]) // PARTITION BY a + b + .order_by(vec![col("c").sort(true, true)]) + .build() + .unwrap(); + + let plan = LogicalPlanBuilder::from(table_scan) + .window(vec![window])? + .filter(in_subquery(add(col("a"), col("b")), subplan))? + .build()?; + assert_plan_not_transformed!(plan.clone()); + + assert_optimized_plan_equal!( + plan, + @r" + Filter: test.a + test.b IN () + Subquery: + Projection: sq.c + TableScan: sq, full_filters=[sq.c = outer_ref(test.a)] + WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] + TableScan: test + " + ) + } + /// An immutable one-argument UDF call, for building an expression partition /// key that is a function call rather than an arithmetic expression. fn immutable_udf_call(arg: Expr) -> Expr { From adb07057b8db1decb7b46f61f46d802134391d97 Mon Sep 17 00:00:00 2001 From: Qi Zhu Date: Wed, 30 Sep 2026 12:13:16 +0800 Subject: [PATCH 5/6] test: cover window filter pushdown over an expression key end to end The unit tests pin where the filter lands. Nothing pinned that the answers stay the same, and the sqllogictest suite had no query at all with a filter above a window partitioned by an expression, which is why this change moved no golden plan in the whole suite. Five cases over one small table, asserting results rather than plans so they do not churn with plan formatting: - the unfiltered baseline, pinning each partition's sum - a predicate on the partition key, where the partition survives whole - a predicate on a non-key column, which is the one a wrong push breaks: pushing it would recompute partition 'a' over the survivors and report 5 instead of 6 - an IN (subquery) predicate, guarding the subquery case - a predicate built on top of the key --- .../push_down_filter_regression.slt | 89 +++++++++++++++++++ 1 file changed, 89 insertions(+) diff --git a/datafusion/sqllogictest/test_files/push_down_filter_regression.slt b/datafusion/sqllogictest/test_files/push_down_filter_regression.slt index 20962e15486de..4ad351767c575 100644 --- a/datafusion/sqllogictest/test_files/push_down_filter_regression.slt +++ b/datafusion/sqllogictest/test_files/push_down_filter_regression.slt @@ -723,3 +723,92 @@ query I SELECT sum(c) FROM (SELECT random() < 0.5 AS k, count(*) AS c FROM generate_series(1, 10000) GROUP BY random() < 0.5) WHERE k OR NOT k; ---- 10000 + +# Window filter pushdown over an expression PARTITION BY key. +# +# A predicate that reads only the partition keys is constant within a partition, +# so pushing it below the window drops whole partitions and leaves every +# surviving row's window value alone. The unit tests pin where the filter lands; +# these pin that the answers do not move, which is what a wrong push breaks. + +statement ok +create table window_expr_key(k varchar, v int) as values + ('a', 1), ('a', 2), ('a', 3), + ('', 4), ('', 5), + ('b', 6); + +# Partitions under NULLIF(k, '') are 'a' => {1,2,3}, NULL => {4,5}, 'b' => {6}. +query TII +SELECT k, v, s FROM ( + SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s + FROM window_expr_key +) +ORDER BY v; +---- +a 1 6 +a 2 6 +a 3 6 +(empty) 4 9 +(empty) 5 9 +b 6 6 + +# Pushable: the predicate reads only the partition key, so partition 'a' +# survives whole and its sum is still 1+2+3. +query TII +SELECT k, v, s FROM ( + SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s + FROM window_expr_key +) +WHERE NULLIF(k, '') = 'a' +ORDER BY v; +---- +a 1 6 +a 2 6 +a 3 6 + +# Not pushable: `v` is not a partition key, so this takes rows out of a +# partition rather than dropping partitions. Pushing it would recompute the +# sums over the survivors and report 5 instead of 6 for partition 'a'. +query TII +SELECT k, v, s FROM ( + SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s + FROM window_expr_key +) +WHERE v > 1 +ORDER BY v; +---- +a 2 6 +a 3 6 +(empty) 4 9 +(empty) 5 9 +b 6 6 + +# Not pushable: a subquery hides what it correlates on from the expression +# walk, so it has to stay above the window even though its left-hand side is +# the partition key. +query TII +SELECT k, v, s FROM ( + SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s + FROM window_expr_key +) +WHERE NULLIF(k, '') IN (SELECT k FROM window_expr_key WHERE v > 5) +ORDER BY v; +---- +b 6 6 + +# A predicate built on top of the key still reads only the key. +query TII +SELECT k, v, s FROM ( + SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s + FROM window_expr_key +) +WHERE NULLIF(k, '') IS NOT NULL +ORDER BY v; +---- +a 1 6 +a 2 6 +a 3 6 +b 6 6 + +statement ok +drop table window_expr_key; From 290995848442d5a8009fc0ba4d08abeb8850ed87 Mon Sep 17 00:00:00 2001 From: Qi Zhu Date: Fri, 9 Oct 2026 14:50:47 +0800 Subject: [PATCH 6/6] Address review: keys by reference, comments with examples, coverage in slt - extract_partition_keys returns HashSet<&Expr>; no clone of the key exprs - comments no longer describe the previous behavior; the helper's doc gives examples of what is and is not pushed under PARTITION BY a, b + c - the DataFrame-built unit tests are replaced by EXPLAIN + result cases in push_down_filter_regression.slt (expression key, key used twice, mixed column and expression keys, operand order, several windows with common and with disjoint keys, subqueries, volatile predicates); one unit test remains for the expression key --- datafusion/optimizer/src/push_down_filter.rs | 530 +----------------- .../push_down_filter_regression.slt | 241 +++++++- 2 files changed, 245 insertions(+), 526 deletions(-) diff --git a/datafusion/optimizer/src/push_down_filter.rs b/datafusion/optimizer/src/push_down_filter.rs index d3dd6f6c3e1f8..efd5a61e4ed9f 100644 --- a/datafusion/optimizer/src/push_down_filter.rs +++ b/datafusion/optimizer/src/push_down_filter.rs @@ -1135,13 +1135,8 @@ impl OptimizerRule for PushDownFilter { // multiple window functions, each with potentially different partition keys. // Therefore, we need to ensure that any potential partition key returned is used in // ALL window functions. Otherwise, filters cannot be pushed by through that column. - // Keyed by the partition *expression*, not by a name synthesised - // from it. `PARTITION BY a + b` used to be mapped through - // `qualified_name()` to a column literally called "a + b", which no - // predicate's real column refs could ever match, so such a key was - // dead weight in this set. - fn extract_partition_keys(func: &WindowFunction) -> HashSet { - func.params.partition_by.iter().cloned().collect() + fn extract_partition_keys(func: &WindowFunction) -> HashSet<&Expr> { + func.params.partition_by.iter().collect() } let potential_partition_keys = window @@ -1177,12 +1172,8 @@ impl OptimizerRule for PushDownFilter { let mut keep_predicates = vec![]; let mut push_predicates = vec![]; for expr in predicates { - // A volatile predicate has to stay above the window. Pushing it - // changes which rows the window function sees, and so the value - // it computes for the rows that do survive. Checking this first - // also covers a volatile predicate that reads no columns at all, - // such as `random() < 0.5`, which would otherwise satisfy the - // partition-key test vacuously. + // A volatile predicate has to stay above the window: pushing it + // changes which rows the window function sees. if !expr.is_volatile() && reads_only_partition_keys(&expr, &potential_partition_keys)? { @@ -1596,33 +1587,27 @@ fn with_filters(predicates: Vec, plan: LogicalPlan) -> LogicalPlan { } } -/// Does `expr` read nothing beyond the given window partition keys? +/// Can `expr` be evaluated below a window with these `PARTITION BY` keys? /// -/// A subtree that is exactly one of the keys counts as read in full, so a -/// predicate on an *expression* key, say `NULLIF(c, '') IS NOT NULL` against -/// `PARTITION BY NULLIF(c, '')`, qualifies even though the column it ultimately -/// reads (`c`) is not a key on its own. Such a predicate is constant within each -/// partition, so applying it below the window drops whole partitions and leaves -/// every surviving row's window value unchanged. +/// A predicate that reads only the partition keys is constant within each +/// partition, so filtering before the window drops whole partitions and leaves +/// the surviving rows' window values unchanged. "Reads only the keys" is checked +/// structurally: every column reference must sit inside a subtree that is equal +/// to one of the keys. Given `PARTITION BY a, b + c`: /// -/// Matching is structural, which makes this conservative rather than wrong: a -/// predicate written `b + a` does not match a key written `a + b`, and is simply -/// left above the window. +/// * `a < 5`, `b + c = 4` and `(b + c) + 1 > 10` can be pushed down +/// * `d < 5` and `b < 5` cannot (`b` on its own is not a key), and neither can +/// `c + b = 4` (the match is structural, `c + b` is not `b + c`) /// -/// A node carrying a subquery counts as reading something else, whatever the -/// keys are. `Expr::apply` does not descend into a subquery's plan, so the -/// columns it correlates on are invisible here, and an outer reference can vary -/// inside a single partition: with `PARTITION BY a + b`, the predicate -/// `a + b > (SELECT ... WHERE inner.x = outer.a)` would otherwise look like it -/// reads nothing but the key. The same blind spot hides a volatile function -/// inside a subquery from the caller's `is_volatile` check. +/// A predicate containing a subquery is never pushed: what the subquery reads is +/// not visible from the expression tree. fn reads_only_partition_keys( expr: &Expr, - partition_keys: &HashSet, + partition_keys: &HashSet<&Expr>, ) -> Result { let mut reads_something_else = false; expr.apply(|node| { - Ok(if partition_keys.contains(node) { + Ok(if partition_keys.contains(&node) { // the whole key was matched, so whatever it reads is accounted for TreeNodeRecursion::Jump } else if reads_beyond_this_node(node) { @@ -1635,12 +1620,9 @@ fn reads_only_partition_keys( Ok(!reads_something_else) } -/// Does this node read data that walking its children cannot account for? -/// -/// A column reads itself. The subquery-bearing variants read whatever their -/// plan reads, including outer references to the window's own input, and -/// `Expr::apply` yields none of that: `Exists` and `ScalarSubquery` are leaves, -/// and `InSubquery` and `SetComparison` expose only their left-hand expression. +/// Does this node read data that walking its children cannot account for? A +/// column reads itself; a subquery reads whatever its plan reads, which +/// `Expr::apply` does not visit. fn reads_beyond_this_node(node: &Expr) -> bool { matches!( node, @@ -1679,7 +1661,6 @@ mod tests { ScalarFunctionArgs, ScalarUDF, ScalarUDFImpl, Signature, Subquery, TableScan, TableSource, TableType, UserDefinedLogicalNodeCore, Volatility, WindowFunctionDefinition, col, exists, in_list, in_subquery, lit, out_ref_col, - scalar_subquery, }; use crate::OptimizerContext; @@ -2025,8 +2006,10 @@ mod tests { ) } - /// verifies that a filter on an expression partition key is pushed, matched - /// against the key expression itself rather than against its column refs + /// verifies that a filter on an expression partition key is pushed; the + /// remaining shapes (mixed keys, operand order, subqueries, volatile + /// predicates, several windows) are covered in + /// `push_down_filter_regression.slt` #[test] fn filter_expression_move_window() -> Result<()> { let table_scan = test_table_scan()?; @@ -2044,9 +2027,7 @@ mod tests { let plan = LogicalPlanBuilder::from(table_scan) .window(vec![window])? - // the single partition column "test.a + test.b" is not available to the - // plan, so the predicate is written over the underlying columns - .filter(add(col("a"), col("b")).gt(lit(10i64)))? + .filter(add(col("a"), col("b")).gt(lit(10i64)))? // a + b > 10 .build()?; assert_optimized_plan_equal!( @@ -2058,467 +2039,6 @@ mod tests { ) } - /// verifies that a single predicate spanning a column key and an expression - /// key is pushed: it is constant within every partition, but neither the old - /// column-name matching nor a rule that looked only at expression keys could - /// see that - #[test] - fn filter_move_window_mixed_column_and_expression_keys() -> Result<()> { - let table_scan = test_table_scan()?; - - let window = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![col("a"), add(col("a"), col("b"))]) // PARTITION BY a, a + b - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let plan = LogicalPlanBuilder::from(table_scan) - .window(vec![window])? - .filter( - col("a") - .gt(lit(1i64)) - .or(add(col("a"), col("b")).gt(lit(10i64))), - )? - .build()?; - - assert_optimized_plan_equal!( - plan, - @r" - WindowAggr: windowExpr=[[rank() PARTITION BY [test.a, test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] - TableScan: test, full_filters=[test.a > Int64(1) OR test.a + test.b > Int64(10)] - " - ) - } - - /// verifies that an expression partition key does not make the columns it - /// reads pushable on their own: `a` alone is not constant within a partition - /// of `a + b` - #[test] - fn filter_keep_window_column_underlying_expression_key() -> Result<()> { - let table_scan = test_table_scan()?; - - let window = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![add(col("a"), col("b"))]) // PARTITION BY a + b - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let plan = LogicalPlanBuilder::from(table_scan) - .window(vec![window])? - .filter(col("a").gt(lit(10i64)))? - .build()?; - assert_plan_not_transformed!(plan.clone()); - - assert_optimized_plan_equal!( - plan, - @r" - Filter: test.a > Int64(10) - WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] - TableScan: test - " - ) - } - - /// verifies that a predicate on a *function* expression key is pushed. This is - /// the shape the rule exists for: `PARTITION BY NULLIF(c, '')` with a predicate - /// on `NULLIF(c, '')`, where the key is not an arithmetic expression and the - /// column it reads is not a key on its own. - #[test] - fn filter_move_window_function_expression_key() -> Result<()> { - let table_scan = test_table_scan()?; - - let window = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![immutable_udf_call(col("c"))]) - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let plan = LogicalPlanBuilder::from(table_scan) - .window(vec![window])? - .filter(immutable_udf_call(col("c")).is_not_null())? - .build()?; - - assert_optimized_plan_equal!( - plan, - @r" - WindowAggr: windowExpr=[[rank() PARTITION BY [TestScalarUDF(test.c)] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] - TableScan: test, full_filters=[TestScalarUDF(test.c) IS NOT NULL] - " - ) - } - - /// verifies that matching is structural, so a predicate written with the - /// operands in the other order does not match the key and stays above the - /// window. Conservative rather than wrong. - #[test] - fn filter_keep_window_expression_key_operand_order() -> Result<()> { - let table_scan = test_table_scan()?; - - let window = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![add(col("a"), col("b"))]) // PARTITION BY a + b - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let plan = LogicalPlanBuilder::from(table_scan) - .window(vec![window])? - .filter(add(col("b"), col("a")).gt(lit(10i64)))? // b + a, not a + b - .build()?; - assert_plan_not_transformed!(plan.clone()); - - assert_optimized_plan_equal!( - plan, - @r" - Filter: test.b + test.a > Int64(10) - WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] - TableScan: test - " - ) - } - - /// verifies that a predicate built *on top of* an expression key is pushed: - /// the key subtree accounts for everything it reads, and the literal around it - /// reads nothing, so the predicate is still constant within each partition - #[test] - fn filter_move_window_expression_key_subtree() -> Result<()> { - let table_scan = test_table_scan()?; - - let window = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![add(col("a"), col("b"))]) // PARTITION BY a + b - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let plan = LogicalPlanBuilder::from(table_scan) - .window(vec![window])? - .filter(add(add(col("a"), col("b")), lit(1i64)).gt(lit(10i64)))? - .build()?; - - assert_optimized_plan_equal!( - plan, - @r" - WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] - TableScan: test, full_filters=[test.a + test.b + Int64(1) > Int64(10)] - " - ) - } - - /// verifies that an expression key shared by every window function is pushed, - /// the expression-key counterpart of filter_multiple_windows_common_partitions - #[test] - fn filter_multiple_windows_common_expression_partition() -> Result<()> { - let table_scan = test_table_scan()?; - - let window1 = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![add(col("a"), col("b"))]) - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let window2 = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![col("b"), add(col("a"), col("b"))]) - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let plan = LogicalPlanBuilder::from(table_scan) - .window(vec![window1, window2])? - .filter(add(col("a"), col("b")).gt(lit(10i64)))? // a + b is in both - .build()?; - - assert_optimized_plan_equal!( - plan, - @r" - WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, rank() PARTITION BY [test.b, test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] - TableScan: test, full_filters=[test.a + test.b > Int64(10)] - " - ) - } - - /// verifies that an expression key present in only one of the window functions - /// is not pushed, since the other function's partitions would change - #[test] - fn filter_multiple_windows_disjoint_expression_partition() -> Result<()> { - let table_scan = test_table_scan()?; - - let window1 = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![add(col("a"), col("b"))]) - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let window2 = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![col("a")]) - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let plan = LogicalPlanBuilder::from(table_scan) - .window(vec![window1, window2])? - .filter(add(col("a"), col("b")).gt(lit(10i64)))? // a + b is in one only - .build()?; - assert_plan_not_transformed!(plan.clone()); - - assert_optimized_plan_equal!( - plan, - @r" - Filter: test.a + test.b > Int64(10) - WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, rank() PARTITION BY [test.a] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] - TableScan: test - " - ) - } - - /// verifies that the volatile check still wins over an expression key: the - /// predicate matches the key structurally, but pushing it would change which - /// rows each partition contains and so the window values of the survivors - #[test] - fn filter_volatile_keep_window_expression_key() -> Result<()> { - let table_scan = test_table_scan()?; - - let window = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![add(col("a"), col("b"))]) - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let fun = ScalarUDF::new_from_impl(TestScalarUDF { - signature: Signature::exact(vec![], Volatility::Volatile), - }); - let volatile = - Expr::ScalarFunction(ScalarFunction::new_udf(Arc::new(fun), vec![])); - - let plan = LogicalPlanBuilder::from(table_scan) - .window(vec![window])? - .filter(add(col("a"), col("b")).gt(volatile))? - .build()?; - assert_plan_not_transformed!(plan.clone()); - - assert_optimized_plan_equal!( - plan, - @r" - Filter: test.a + test.b > TestScalarUDF() - WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] - TableScan: test - " - ) - } - - /// verifies that a predicate carrying a correlated scalar subquery stays - /// above the window even though its visible part is exactly the expression - /// key. `Expr::apply` does not enter the subquery plan, so `outer.a`, which - /// varies inside a partition of `a + b`, is invisible to the key test. - #[test] - fn filter_keep_window_scalar_subquery_over_expression_key() -> Result<()> { - let table_scan = test_table_scan()?; - let subplan = Arc::new( - LogicalPlanBuilder::from(test_table_scan_with_name("sq")?) - .filter(col("sq.c").eq(out_ref_col(DataType::UInt32, "test.a")))? - .aggregate(Vec::::new(), vec![sum(col("sq.c"))])? - .build()?, - ); - - let window = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![add(col("a"), col("b"))]) // PARTITION BY a + b - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let plan = LogicalPlanBuilder::from(table_scan) - .window(vec![window])? - .filter(add(col("a"), col("b")).gt(scalar_subquery(subplan)))? - .build()?; - assert_plan_not_transformed!(plan.clone()); - - assert_optimized_plan_equal!( - plan, - @r" - Filter: test.a + test.b > () - Subquery: - Aggregate: groupBy=[[]], aggr=[[sum(sq.c)]] - TableScan: sq, full_filters=[sq.c = outer_ref(test.a)] - WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] - TableScan: test - " - ) - } - - /// the same guard for `IN (subquery)`, whose subquery plan `Expr::apply` - /// also skips: only the left-hand expression is walked - #[test] - fn filter_keep_window_in_subquery_over_expression_key() -> Result<()> { - let table_scan = test_table_scan()?; - let subplan = Arc::new( - LogicalPlanBuilder::from(test_table_scan_with_name("sq")?) - .filter(col("sq.c").eq(out_ref_col(DataType::UInt32, "test.a")))? - .project(vec![col("sq.c")])? - .build()?, - ); - - let window = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![add(col("a"), col("b"))]) // PARTITION BY a + b - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let plan = LogicalPlanBuilder::from(table_scan) - .window(vec![window])? - .filter(in_subquery(add(col("a"), col("b")), subplan))? - .build()?; - assert_plan_not_transformed!(plan.clone()); - - assert_optimized_plan_equal!( - plan, - @r" - Filter: test.a + test.b IN () - Subquery: - Projection: sq.c - TableScan: sq, full_filters=[sq.c = outer_ref(test.a)] - WindowAggr: windowExpr=[[rank() PARTITION BY [test.a + test.b] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] - TableScan: test - " - ) - } - - /// An immutable one-argument UDF call, for building an expression partition - /// key that is a function call rather than an arithmetic expression. - fn immutable_udf_call(arg: Expr) -> Expr { - let fun = ScalarUDF::new_from_impl(TestScalarUDF { - signature: Signature::any(1, Volatility::Immutable), - }); - Expr::ScalarFunction(ScalarFunction::new_udf(Arc::new(fun), vec![arg])) - } - - /// verifies that a volatile predicate is not pushed through a window, even - /// when it reads no columns and so trivially satisfies the partition-key test - #[test] - fn filter_volatile_keep_window() -> Result<()> { - let table_scan = test_table_scan()?; - - let window = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![col("a")]) - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let fun = ScalarUDF::new_from_impl(TestScalarUDF { - signature: Signature::exact(vec![], Volatility::Volatile), - }); - let volatile = - Expr::ScalarFunction(ScalarFunction::new_udf(Arc::new(fun), vec![])); - - let plan = LogicalPlanBuilder::from(table_scan) - .window(vec![window])? - .filter(volatile.gt(lit(10i64)))? - .build()?; - - assert_optimized_plan_equal!( - plan, - @r" - Filter: TestScalarUDF() > Int64(10) - WindowAggr: windowExpr=[[rank() PARTITION BY [test.a] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] - TableScan: test - " - ) - } - - /// verifies that filters are not pushed on order by columns (that are not used in partitioning) - #[test] - fn filter_order_keep_window() -> Result<()> { - let table_scan = test_table_scan()?; - - let window = Expr::from(WindowFunction::new( - WindowFunctionDefinition::WindowUDF( - datafusion_functions_window::rank::rank_udwf(), - ), - vec![], - )) - .partition_by(vec![col("a")]) - .order_by(vec![col("c").sort(true, true)]) - .build() - .unwrap(); - - let plan = LogicalPlanBuilder::from(table_scan) - .window(vec![window])? - .filter(col("c").gt(lit(10i64)))? - .build()?; - assert_plan_not_transformed!(plan.clone()); - - assert_optimized_plan_equal!( - plan, - @r" - Filter: test.c > Int64(10) - WindowAggr: windowExpr=[[rank() PARTITION BY [test.a] ORDER BY [test.c ASC NULLS FIRST] ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] - TableScan: test - " - ) - } - /// verifies that when we use multiple window functions with a common partition key, the filter /// on that key is pushed #[test] diff --git a/datafusion/sqllogictest/test_files/push_down_filter_regression.slt b/datafusion/sqllogictest/test_files/push_down_filter_regression.slt index 4ad351767c575..f85d568efc789 100644 --- a/datafusion/sqllogictest/test_files/push_down_filter_regression.slt +++ b/datafusion/sqllogictest/test_files/push_down_filter_regression.slt @@ -724,12 +724,13 @@ SELECT sum(c) FROM (SELECT random() < 0.5 AS k, count(*) AS c FROM generate_seri ---- 10000 -# Window filter pushdown over an expression PARTITION BY key. +# Window filter pushdown over expression PARTITION BY keys. # # A predicate that reads only the partition keys is constant within a partition, -# so pushing it below the window drops whole partitions and leaves every -# surviving row's window value alone. The unit tests pin where the filter lands; -# these pin that the answers do not move, which is what a wrong push breaks. +# so it is pushed below the window; anything else stays above it. + +statement ok +set datafusion.explain.logical_plan_only = true; statement ok create table window_expr_key(k varchar, v int) as values @@ -752,8 +753,20 @@ a 3 6 (empty) 5 9 b 6 6 -# Pushable: the predicate reads only the partition key, so partition 'a' -# survives whole and its sum is still 1+2+3. +# Pushed: the predicate reads only the partition key. +query TT +EXPLAIN SELECT k, v, s FROM ( + SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s + FROM window_expr_key +) +WHERE NULLIF(k, '') = 'a'; +---- +logical_plan +01)Projection: window_expr_key.k, window_expr_key.v, sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS s +02)--WindowAggr: windowExpr=[[sum(CAST(window_expr_key.v AS Int64)) PARTITION BY [nullif(window_expr_key.k, Utf8View(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING]] +03)----Filter: nullif(window_expr_key.k, Utf8View("")) = Utf8View("a") +04)------TableScan: window_expr_key projection=[k, v] + query TII SELECT k, v, s FROM ( SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s @@ -766,49 +779,235 @@ a 1 6 a 2 6 a 3 6 -# Not pushable: `v` is not a partition key, so this takes rows out of a -# partition rather than dropping partitions. Pushing it would recompute the -# sums over the survivors and report 5 instead of 6 for partition 'a'. +# Pushed: a predicate built on top of the key still reads only the key. +query TT +EXPLAIN SELECT k, v, s FROM ( + SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s + FROM window_expr_key +) +WHERE NULLIF(k, '') IS NOT NULL; +---- +logical_plan +01)Projection: window_expr_key.k, window_expr_key.v, sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS s +02)--WindowAggr: windowExpr=[[sum(CAST(window_expr_key.v AS Int64)) PARTITION BY [nullif(window_expr_key.k, Utf8View(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING]] +03)----Filter: nullif(window_expr_key.k, Utf8View("")) IS NOT NULL +04)------TableScan: window_expr_key projection=[k, v] + query TII SELECT k, v, s FROM ( SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s FROM window_expr_key ) -WHERE v > 1 +WHERE NULLIF(k, '') IS NOT NULL ORDER BY v; ---- +a 1 6 a 2 6 a 3 6 -(empty) 4 9 -(empty) 5 9 b 6 6 -# Not pushable: a subquery hides what it correlates on from the expression -# walk, so it has to stay above the window even though its left-hand side is -# the partition key. +# Kept: `v` is not a partition key. Pushing it would recompute the sums over the +# surviving rows and report 5 instead of 6 for partition 'a'. +query TT +EXPLAIN SELECT k, v, s FROM ( + SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s + FROM window_expr_key +) +WHERE v > 1; +---- +logical_plan +01)Projection: window_expr_key.k, window_expr_key.v, sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS s +02)--Filter: window_expr_key.v > Int32(1) +03)----WindowAggr: windowExpr=[[sum(CAST(window_expr_key.v AS Int64)) PARTITION BY [nullif(window_expr_key.k, Utf8View(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING]] +04)------TableScan: window_expr_key projection=[k, v] + query TII SELECT k, v, s FROM ( SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s FROM window_expr_key ) -WHERE NULLIF(k, '') IN (SELECT k FROM window_expr_key WHERE v > 5) +WHERE v > 1 ORDER BY v; ---- +a 2 6 +a 3 6 +(empty) 4 9 +(empty) 5 9 b 6 6 -# A predicate built on top of the key still reads only the key. +# Kept: `k` on its own is not a key, only NULLIF(k, '') is. +query TT +EXPLAIN SELECT k, v, s FROM ( + SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s + FROM window_expr_key +) +WHERE k = 'a'; +---- +logical_plan +01)Projection: window_expr_key.k, window_expr_key.v, sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS s +02)--Filter: window_expr_key.k = Utf8View("a") +03)----WindowAggr: windowExpr=[[sum(CAST(window_expr_key.v AS Int64)) PARTITION BY [nullif(window_expr_key.k, Utf8View(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING]] +04)------TableScan: window_expr_key projection=[k, v] + +# Kept: a subquery, even one whose left-hand side is the key. +query TT +EXPLAIN SELECT k, v, s FROM ( + SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s + FROM window_expr_key +) +WHERE NULLIF(k, '') IN (SELECT k FROM window_expr_key WHERE v > 5); +---- +logical_plan +01)LeftSemi Join: nullif(window_expr_key.k, Utf8View("")) = __correlated_sq_1.k +02)--Projection: window_expr_key.k, window_expr_key.v, sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS s +03)----WindowAggr: windowExpr=[[sum(CAST(window_expr_key.v AS Int64)) PARTITION BY [nullif(window_expr_key.k, Utf8View(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING]] +04)------TableScan: window_expr_key projection=[k, v] +05)--SubqueryAlias: __correlated_sq_1 +06)----Projection: window_expr_key.k +07)------Filter: window_expr_key.v > Int32(5) +08)--------TableScan: window_expr_key projection=[k, v] + query TII SELECT k, v, s FROM ( SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s FROM window_expr_key ) -WHERE NULLIF(k, '') IS NOT NULL +WHERE NULLIF(k, '') IN (SELECT k FROM window_expr_key WHERE v > 5) ORDER BY v; ---- -a 1 6 -a 2 6 -a 3 6 b 6 6 +# Kept: a scalar subquery. +query TT +EXPLAIN SELECT k, v, s FROM ( + SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s + FROM window_expr_key +) +WHERE NULLIF(k, '') = (SELECT max(k) FROM window_expr_key); +---- +logical_plan +01)Projection: window_expr_key.k, window_expr_key.v, sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS s +02)--Filter: nullif(window_expr_key.k, Utf8View("")) = () +03)----Subquery: +04)------Aggregate: groupBy=[[]], aggr=[[max(window_expr_key.k)]] +05)--------TableScan: window_expr_key projection=[k] +06)----WindowAggr: windowExpr=[[sum(CAST(window_expr_key.v AS Int64)) PARTITION BY [nullif(window_expr_key.k, Utf8View(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING]] +07)------TableScan: window_expr_key projection=[k, v] + +# Kept: a volatile predicate, even one that reads no column at all. +query TT +EXPLAIN SELECT k, v, s FROM ( + SELECT k, v, SUM(v) OVER (PARTITION BY NULLIF(k, '')) AS s + FROM window_expr_key +) +WHERE random() < 2; +---- +logical_plan +01)Projection: window_expr_key.k, window_expr_key.v, sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS s +02)--Filter: random() < Float64(2) +03)----WindowAggr: windowExpr=[[sum(CAST(window_expr_key.v AS Int64)) PARTITION BY [nullif(window_expr_key.k, Utf8View(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS sum(window_expr_key.v) PARTITION BY [nullif(window_expr_key.k,Utf8(""))] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING]] +04)------TableScan: window_expr_key projection=[k, v] + statement ok drop table window_expr_key; + +statement ok +create table window_ab(a int, b int, c int) as values + (1, 1, 10), (1, 2, 20), (2, 2, 30), (3, 3, 40); + +# Pushed: an arithmetic expression key. +query TT +EXPLAIN SELECT a, b, s FROM ( + SELECT a, b, SUM(c) OVER (PARTITION BY a + b) AS s FROM window_ab +) +WHERE a + b > 2; +---- +logical_plan +01)Projection: window_ab.a, window_ab.b, sum(window_ab.c) PARTITION BY [window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS s +02)--WindowAggr: windowExpr=[[sum(CAST(window_ab.c AS Int64)) PARTITION BY [window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING]] +03)----Filter: window_ab.a + window_ab.b > Int32(2) +04)------TableScan: window_ab projection=[a, b, c] + +# Pushed: the key appears twice in one predicate. +query TT +EXPLAIN SELECT a, b, s FROM ( + SELECT a, b, SUM(c) OVER (PARTITION BY a + b) AS s FROM window_ab +) +WHERE a + b > 5 OR a + b < 3; +---- +logical_plan +01)Projection: window_ab.a, window_ab.b, sum(window_ab.c) PARTITION BY [window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS s +02)--WindowAggr: windowExpr=[[sum(CAST(window_ab.c AS Int64)) PARTITION BY [window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING]] +03)----Projection: window_ab.a, window_ab.b, window_ab.c +04)------Filter: __common_expr_3 > Int32(5) OR __common_expr_3 < Int32(3) +05)--------Projection: window_ab.a + window_ab.b AS __common_expr_3, window_ab.a, window_ab.b, window_ab.c +06)----------TableScan: window_ab projection=[a, b, c] + +# Pushed: a predicate over a column key and an expression key of the same window. +query TT +EXPLAIN SELECT a, b, s FROM ( + SELECT a, b, SUM(c) OVER (PARTITION BY a, a + b) AS s FROM window_ab +) +WHERE a > 1 OR a + b > 4; +---- +logical_plan +01)Projection: window_ab.a, window_ab.b, sum(window_ab.c) PARTITION BY [window_ab.a, window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS s +02)--WindowAggr: windowExpr=[[sum(CAST(window_ab.c AS Int64)) PARTITION BY [window_ab.a, window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING]] +03)----Filter: window_ab.a > Int32(1) OR window_ab.a + window_ab.b > Int32(4) +04)------TableScan: window_ab projection=[a, b, c] + +# Kept: the match is structural, `b + a` is not the key `a + b`. +query TT +EXPLAIN SELECT a, b, s FROM ( + SELECT a, b, SUM(c) OVER (PARTITION BY a + b) AS s FROM window_ab +) +WHERE b + a > 2; +---- +logical_plan +01)Projection: window_ab.a, window_ab.b, sum(window_ab.c) PARTITION BY [window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS s +02)--Filter: window_ab.b + window_ab.a > Int32(2) +03)----Projection: window_ab.a, window_ab.b, sum(window_ab.c) PARTITION BY [window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING +04)------WindowAggr: windowExpr=[[sum(CAST(window_ab.c AS Int64)) PARTITION BY [window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING]] +05)--------TableScan: window_ab projection=[a, b, c] + +# Pushed: every window function partitions by the key. +query TT +EXPLAIN SELECT a, b, s, r FROM ( + SELECT a, b, + SUM(c) OVER (PARTITION BY a + b) AS s, + RANK() OVER (PARTITION BY a + b ORDER BY c) AS r + FROM window_ab +) +WHERE a + b > 2; +---- +logical_plan +01)Projection: window_ab.a, window_ab.b, sum(window_ab.c) PARTITION BY [window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS s, rank() PARTITION BY [window_ab.a + window_ab.b] ORDER BY [window_ab.c ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS r +02)--WindowAggr: windowExpr=[[sum(CAST(window_ab.c AS Int64)) PARTITION BY [__common_expr_1 AS window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING]] +03)----WindowAggr: windowExpr=[[rank() PARTITION BY [__common_expr_1 AS window_ab.a + window_ab.b] ORDER BY [window_ab.c ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] +04)------Projection: window_ab.a + window_ab.b AS __common_expr_1, window_ab.a, window_ab.b, window_ab.c +05)--------Filter: window_ab.a + window_ab.b > Int32(2) +06)----------TableScan: window_ab projection=[a, b, c] + +# Kept: one window function partitions by something else. +query TT +EXPLAIN SELECT a, b, s, r FROM ( + SELECT a, b, + SUM(c) OVER (PARTITION BY a + b) AS s, + RANK() OVER (PARTITION BY a ORDER BY c) AS r + FROM window_ab +) +WHERE a + b > 2; +---- +logical_plan +01)Projection: window_ab.a, window_ab.b, sum(window_ab.c) PARTITION BY [window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING AS s, rank() PARTITION BY [window_ab.a] ORDER BY [window_ab.c ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS r +02)--Filter: window_ab.a + window_ab.b > Int32(2) +03)----Projection: window_ab.a, window_ab.b, sum(window_ab.c) PARTITION BY [window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING, rank() PARTITION BY [window_ab.a] ORDER BY [window_ab.c ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW +04)------WindowAggr: windowExpr=[[rank() PARTITION BY [window_ab.a] ORDER BY [window_ab.c ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]] +05)--------WindowAggr: windowExpr=[[sum(CAST(window_ab.c AS Int64)) PARTITION BY [window_ab.a + window_ab.b] ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING]] +06)----------TableScan: window_ab projection=[a, b, c] + +statement ok +drop table window_ab; + +statement ok +set datafusion.explain.logical_plan_only = false;