Skip to content

push_down_filter does not push window predicates on expression PARTITION BY keys #25864

Description

@zhuqi-lucas

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:

// 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.
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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

enhancementNew feature or request

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions