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]

Reply via email to