Repository navigation
feat: push window predicates on expression PARTITION BY keys - #25865
zhuqi-lucas wants to merge 7 commits into
Conversation
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.
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.
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #25865 +/- ##
==========================================
- Coverage 82.73% 82.72% -0.01%
==========================================
Files 1147 1147
Lines 449213 449541 +328
Branches 449213 449541 +328
==========================================
+ Hits 371634 371903 +269
- Misses 54929 54932 +3
- Partials 22650 22706 +56 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
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
There was a problem hiding this comment.
Copilot review overview
🔵 Needs a closer look
A subquery predicate may still be pushed unsafely, and the SQL plan changes remain unverified.
Review effort: Balanced
Findings: 1
Open (1)
What changed in this PR
This PR extends DataFusion’s filter pushdown rule to predicates on window expression partition keys and keeps volatile predicates above windows.
Changes:
- Match predicates against partition-key expressions, including mixed column and expression keys.
- Add a volatility guard and unit tests for pushdown and keep-above cases.
| File | Description |
|---|---|
datafusion/optimizer/src/push_down_filter.rs |
Updates window predicate matching and adds regression tests. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
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.
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
alamb
left a comment
There was a problem hiding this comment.
Thank you @zhuqi-lucas -- I think this code and pushdown is correct and makes sense -- thank you for the work
I think we could significantly improve some of the comments and tests, however. I left a bunch of comments
| .filter(add(col("a"), col("b")).gt(lit(10i64)))? | ||
| .build()?; | ||
|
|
||
| assert_optimized_plan_equal!( |
There was a problem hiding this comment.
Is there a reason we can't use .slt for these tests? The setup code for buulding the dataframe is quite verbose compare to the same thing as SQL, and I found it hard to map from the dataframe API expressions to the SQL level expressions (e.g. add(col(a), col(b)) in the setup code, but a + b in the explain)
I am wondering if there is some way to reduce the size of this PR
| // ALL window functions. Otherwise, filters cannot be pushed by through that column. | ||
| fn extract_partition_keys(func: &WindowFunction) -> HashSet<Column> { | ||
| expr_columns(&func.params.partition_by) | ||
| // Keyed by the partition *expression*, not by a name synthesised |
There was a problem hiding this comment.
Is it important to explain here what used to happen? This comment on the old behavior seems like it will be irrelevant once the PR will merge (it would be better as a comment on the PR I think, not in the code)
| 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 |
There was a problem hiding this comment.
I am not sure "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."
is needed as it explain what the code does, doesn't it?
|
|
||
| /// 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 |
There was a problem hiding this comment.
I found "is exactly one of the keys counts as read in full" very hard to understand
I think this is trying to explain that we can pass any expression down that only refers to partition columns (as in is an expression of columns that only appear in the expression)?
There was a problem hiding this comment.
I think this comment would be much more helpful with some examples
Something like
given PARTITION BY (a, b)
Can push down filters: a < 5, a+b = 4 etc
Can not push down filters like c < 5
| /// 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 |
There was a problem hiding this comment.
I found this paragraph very hard to understand -- I don't think we need to explain the intricate details of why subqueries can't be pushed down
| ('', 4), ('', 5), | ||
| ('b', 6); | ||
|
|
||
| # Partitions under NULLIF(k, '') are 'a' => {1,2,3}, NULL => {4,5}, 'b' => {6}. |
There was a problem hiding this comment.
can you also please add EXPLAIN to these tests so we can see the shape?
| } | ||
|
|
||
| /// 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 |
There was a problem hiding this comment.
again, a reference to the old code is not helpful in code comments -- it is helpful in the context of a PR and I think should be a comment on the PR
|
|
||
| let plan = LogicalPlanBuilder::from(table_scan) | ||
| .window(vec![window])? | ||
| .filter(add(add(col("a"), col("b")), lit(1i64)).gt(lit(10i64)))? |
There was a problem hiding this comment.
would help here to note with ((a + b) + 1) > 10
| ) | ||
| } | ||
|
|
||
| /// verifies that an expression key shared by every window function is pushed, |
There was a problem hiding this comment.
I also recommend a test that has multiple instances of expressions in partition by
PARTITION BY (a+b)
...
WHERE ((a +b) > 5 OR (a + b) < 10)| // 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<Expr> { | ||
| func.params.partition_by.iter().cloned().collect() |
There was a problem hiding this comment.
do you need to clone here? Can we just use HashSet<&Expr> to keep references rather tahn deeply clone the exprs?

Which issue does this PR close?
Rationale for this change
A predicate that depends only on a window's
PARTITION BYkeys is constant within every partition, so applying it below the window drops whole partitions and leaves every surviving row's window value unchanged.push_down_filteralready does this, but only for plain column keys.The key set is built by mapping each partition key through
qualified_name()into aColumn, soPARTITION BY a + bbecomes a column literally named"a + b". A predicate ona + breads the real columnsaandb, which never match that synthesised name, so it stays above the window. The optimization simply does not exist for expression keys such asa + b,NULLIF(c, '')orCOALESCE(x, y).It also affects mixed predicates. With
PARTITION BY year, num * num:the predicate is constant within every partition and safe to push, but its column refs are
{year, num}against a key set of{year, "num * num"}, sonumis not found and the whole conjunct is kept.What changes are included in this PR?
Keep a volatile predicate above the window (first commit, independent of the rest)
A volatile predicate that reads no columns, such as
random() < 0.5, satisfies the subset test vacuously and was pushed below the window, where it changes which rows the window function sees. The aggregate arm already drops volatile group expressions before the equivalent check; this does the same here. Split out as its own commit so it can be taken separately.Match against the key expressions (second commit)
The key set holds the partition key
Exprs rather than names synthesised from them, and each conjunct is walked: a subtree exactly equal to one of the keys counts as read in full, and the predicate is rejected only on reaching aColumnthat no key covered.This is a strict generalization. For a plain column key, the reference to that column is itself a subtree equal to the key, so every predicate pushed today is still pushed; expression keys and mixed predicates are added on top. Matching is structural, so it stays conservative: a predicate written
b + adoes not match a key writtena + band is left in place.Keep a subquery predicate above the window (third commit, from Copilot's review of this PR)
Expr::applydoes not enter a subquery's plan, so the columns it correlates on are invisible to the key test:ExistsandScalarSubqueryare leaves, andInSubqueryandSetComparisonexpose only their left-hand expression. WithPARTITION BY a + b, the predicatelooks like it reads nothing but the key.
outer.avaries 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 theis_volatile()check that runs first.The previous column-key matching happened to keep this expression above the window, because the predicate's column refs
{a, b}never matched the synthesised key name. For a plain column key the hole predates this branch, but expression-key matching would have widened it, so a subquery-bearing node now counts as reading something the walk cannot account for.No predicate rewriting is involved. The existing comment in that arm notes that a window partition expression, unlike an aggregate group expression, is not exposed as a standalone column, so there is nothing to rewrite
a + binto. The predicate is pushed unchanged.Are these changes tested?
Yes, in
push_down_filter's unit tests.Behavior the change turns on:
filter_expression_move_window: a predicate onPARTITION BY a + bis now pushed. This replacesfilter_expression_keep_window, which pinned the previous behavior.filter_move_window_mixed_column_and_expression_keys:PARTITION BY a, a + bwitha > 1 OR a + b > 10is pushed, the case neither the old matching nor an expression-key-only rule could see.filter_move_window_function_expression_key: a function key,PARTITION BY f(c)withf(c) IS NOT NULL, the shape this exists for.filter_move_window_expression_key_subtree: a predicate built on top of a key,(a + b) + 1 > 10, where the key subtree accounts for everything read and the surrounding literal reads nothing.filter_multiple_windows_common_expression_partition: an expression key shared by every window function is pushed.Behavior it deliberately leaves alone, each verified to fail without the corresponding guard:
filter_keep_window_column_underlying_expression_key: an expression key does not make the columns it reads pushable on their own, soa > 10underPARTITION BY a + bstays above.filter_keep_window_expression_key_operand_order: matching is structural, sob + adoes not match a key writtena + b.filter_multiple_windows_disjoint_expression_partition: an expression key present in only one of the window functions is not pushed.filter_volatile_keep_windowandfilter_volatile_keep_window_expression_key: the volatility check still runs first, including when the predicate's operand matches a key.filter_keep_window_scalar_subquery_over_expression_keyandfilter_keep_window_in_subquery_over_expression_key: the subquery guard from the third commit.The existing window tests (
filter_move_window,filter_move_partial_window,filter_order_keep_window,filter_multiple_windows_common_partitions,filter_multiple_windows_disjoint_partitions, and the rest) are unchanged and pass, which is the evidence for the "strict generalization" claim.On plan churn: there is none to report. Matching against the key expressions only ever adds pushes that the name-based matching could not see, and every window plan the test suite pins is either unaffected or already covered by the unit tests above. No
.sltexpectations needed updating.Are there any user-facing changes?
Better plans for window queries filtered on expression partition keys, and a correctness fix for volatile predicates over windows. No API changes.