zhuqi-lucas opened a new issue, #25864:
URL: https://github.com/apache/datafusion/issues/25864
### Is your feature request related to a problem or challenge?
`push_down_filter` pushes a predicate below a `Window` when the predicate
depends only on the window's `PARTITION BY` keys: such a predicate is constant
within every partition, so applying it below the window drops whole partitions
and leaves every surviving row's window value unchanged.
That works for plain column keys, but never fires for *expression* keys such
as `PARTITION BY a + b`, `PARTITION BY NULLIF(c, '')` or `PARTITION BY
COALESCE(x, y)`.
The reason is how the key set is built. Each partition key is mapped through
`qualified_name()` into a `Column`:
```rust
// datafusion/optimizer/src/push_down_filter.rs
fn extract_partition_keys(func: &WindowFunction) -> HashSet<Column> {
expr_columns(&func.params.partition_by)
}
```
so `a + b` becomes a column literally named `"a + b"`. Each conjunct is then
accepted when `expr.column_refs()` is a subset of that set. A predicate on `a +
b` reads the real columns `a` and `b`, which never match the synthesised name,
so the predicate stays above the window.
Two consequences:
1. An expression key can never be matched, so the optimization simply does
not exist for those windows.
2. An expression key also poisons *mixed* predicates. With `PARTITION BY
year, num * num`, the predicate `year = '2021' OR num * num > 4` is constant
within every partition and safe to push, but its column refs are `{year, num}`
while the key set is `{year, "num * num"}`, so `num` is not found and the whole
conjunct is kept above the window.
```sql
SELECT * FROM (
SELECT year, num, SUM(num) OVER (PARTITION BY year, num * num) AS s FROM
t
)
WHERE year = '2021' OR num * num > 4
```
Related, in the same arm: a volatile predicate that reads no columns, such
as `random() < 0.5`, satisfies the subset test vacuously and is pushed below
the window today, where it changes which rows the window function sees. The
aggregate arm already drops volatile group expressions before the equivalent
check.
### Describe the solution you'd like
Match each conjunct against the partition key *expressions* rather than
against names synthesised from them: walk the predicate, treat a subtree that
is exactly one of the keys as satisfied in full, and reject only when a
`Column` is reached that 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 that is
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 +
a` will not match a key written `a + b` and is simply left in place.
### Describe alternatives you've considered
Rewriting the predicate to refer to a synthetic partition column, the way
the aggregate arm does. That does not apply here, because a window partition
expression is not exposed as a standalone column, which is exactly what the
existing comment in that arm points out. No rewriting is needed: the predicate
can be pushed unchanged.
### Additional context
Found while removing a downstream copy of this rule. Atlas carries a
`PushDownWindowPartitionFilter` optimizer rule that does the subtree matching
described above, purely to cover the expression-key case, and it can be deleted
once DataFusion handles it.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]