comphead commented on code in PR #25838:
URL: https://github.com/apache/datafusion/pull/25838#discussion_r4139627021
##########
datafusion/optimizer/src/decorrelate.rs:
##########
@@ -170,18 +176,62 @@ impl TreeNodeRewriter for PullUpCorrelatedExpr {
match (self.exists_sub_query, plan_hold_outer) {
(false, true) => {
// the unsupported case
- self.can_pull_up = false;
- Ok(Transformed::new(plan, false,
TreeNodeRecursion::Jump))
+ Ok(self.stop_pull_up(plan))
}
_ => Ok(Transformed::no(plan)),
}
}
- _ if plan.contains_outer_reference() => {
- // the unsupported cases, the plan expressions contain out
reference columns(like window expressions)
- self.can_pull_up = false;
- Ok(Transformed::new(plan, false, TreeNodeRecursion::Jump))
+ // A correlated filter below these nodes can move above them.
`f_up`
+ // adds the columns the pulled up filter needs to a Projection or
an
+ // Aggregate, and checks that the filter can move above an
+ // Aggregate. A `DISTINCT` gives the same rows when a filter runs
+ // before it or after it, and its output takes the columns the node
Review Comment:
#25808 is a case where this does not hold. The pull up adds the filter's
column to the `DISTINCT` keys, so a `LATERAL` with a non-equality filter
returns duplicates. Worth saying here that `DISTINCT` shares the `Aggregate`
exemption and linking #25808, as the PR description does.
##########
datafusion/optimizer/src/decorrelate.rs:
##########
@@ -170,18 +176,62 @@ impl TreeNodeRewriter for PullUpCorrelatedExpr {
match (self.exists_sub_query, plan_hold_outer) {
(false, true) => {
// the unsupported case
- self.can_pull_up = false;
- Ok(Transformed::new(plan, false,
TreeNodeRecursion::Jump))
+ Ok(self.stop_pull_up(plan))
}
_ => Ok(Transformed::no(plan)),
}
}
- _ if plan.contains_outer_reference() => {
- // the unsupported cases, the plan expressions contain out
reference columns(like window expressions)
- self.can_pull_up = false;
- Ok(Transformed::new(plan, false, TreeNodeRecursion::Jump))
+ // A correlated filter below these nodes can move above them.
`f_up`
+ // adds the columns the pulled up filter needs to a Projection or
an
+ // Aggregate, and checks that the filter can move above an
+ // Aggregate. A `DISTINCT` gives the same rows when a filter runs
+ // before it or after it, and its output takes the columns the node
+ // below it adds. The node itself must not hold an outer reference,
+ // because only a Filter gives its outer references to the join.
+ LogicalPlan::Projection(_)
+ | LogicalPlan::Aggregate(_)
+ | LogicalPlan::Distinct(Distinct::All(_))
+ | LogicalPlan::Join(_)
+ | LogicalPlan::AsOfJoin(_)
+ | LogicalPlan::Repartition(_)
+ | LogicalPlan::SubqueryAlias(_)
+ | LogicalPlan::TableScan(_)
+ | LogicalPlan::EmptyRelation(_)
+ | LogicalPlan::Values(_) => {
+ if plan.contains_outer_reference() {
+ // the unsupported cases, the plan expressions contain out
+ // reference columns
+ Ok(self.stop_pull_up(plan))
+ } else {
+ Ok(Transformed::no(plan))
+ }
+ }
+ // A correlated filter below these nodes cannot move above them:
+ // - A Window function reads the rows of its partition, so the
+ // filter changes the rows each function sees.
+ // - `DISTINCT ON` keeps the first row of each group, so the filter
+ // changes which row is first.
+ // - The pull up does not carry the columns of the pulled up filter
+ // through an Unnest, and fails with a schema error.
+ // - A RecursiveQuery repeats its recursive term.
+ // - The other nodes are not queries and are not in a subquery.
+ LogicalPlan::Window(_)
Review Comment:
This also blocks the case #25792 calls safe. When the correlated filter
reads only `PARTITION BY` columns, it keeps or drops whole partitions, so
pulling it above the window does not change any window value:
```sql
CREATE TABLE o(k INT) AS VALUES (1), (2);
CREATE TABLE i(k INT, v INT) AS VALUES (1, 10), (1, 20), (2, 30);
SELECT o.k, s.v, s.rn FROM o, LATERAL (
SELECT i.v, row_number() OVER (PARTITION BY i.k ORDER BY i.v) AS rn
FROM i WHERE i.k = o.k
) AS s ORDER BY o.k, s.v;
-- expected: (1, 10, 1), (1, 20, 2), (2, 30, 1)
```
If I read `f_down` right, `main` decorrelates this correctly and this PR
turns it into the not-implemented error. No existing test covers it. The check
proposed in the #25792 thread would keep it working: pass `Window` through in
`f_down`, then in `f_up` rebuild it and set `can_pull_up = false` unless every
column from `collect_local_correlated_cols` is a plain `PARTITION BY` column of
every window expression. That is the same continue-then-drop pattern as the
grouping-set check. If you prefer to keep this PR minimal, a test that pins the
error for this query and a follow-up issue would at least record the change.
Longer term, Spark (`DecorrelateInnerQuery`, `Window` case) and DuckDB
(`FlattenDependentJoins::PushDownWindow`) append the correlated columns to
`PARTITION BY`, which would decorrelate #25792 itself instead of erroring.
##########
datafusion/optimizer/src/decorrelate.rs:
##########
@@ -170,18 +176,62 @@ impl TreeNodeRewriter for PullUpCorrelatedExpr {
match (self.exists_sub_query, plan_hold_outer) {
(false, true) => {
// the unsupported case
- self.can_pull_up = false;
- Ok(Transformed::new(plan, false,
TreeNodeRecursion::Jump))
+ Ok(self.stop_pull_up(plan))
}
_ => Ok(Transformed::no(plan)),
}
}
- _ if plan.contains_outer_reference() => {
- // the unsupported cases, the plan expressions contain out
reference columns(like window expressions)
- self.can_pull_up = false;
- Ok(Transformed::new(plan, false, TreeNodeRecursion::Jump))
+ // A correlated filter below these nodes can move above them.
`f_up`
+ // adds the columns the pulled up filter needs to a Projection or
an
+ // Aggregate, and checks that the filter can move above an
+ // Aggregate. A `DISTINCT` gives the same rows when a filter runs
+ // before it or after it, and its output takes the columns the node
+ // below it adds. The node itself must not hold an outer reference,
+ // because only a Filter gives its outer references to the join.
+ LogicalPlan::Projection(_)
+ | LogicalPlan::Aggregate(_)
+ | LogicalPlan::Distinct(Distinct::All(_))
+ | LogicalPlan::Join(_)
+ | LogicalPlan::AsOfJoin(_)
+ | LogicalPlan::Repartition(_)
+ | LogicalPlan::SubqueryAlias(_)
+ | LogicalPlan::TableScan(_)
+ | LogicalPlan::EmptyRelation(_)
+ | LogicalPlan::Values(_) => {
+ if plan.contains_outer_reference() {
+ // the unsupported cases, the plan expressions contain out
+ // reference columns
+ Ok(self.stop_pull_up(plan))
+ } else {
+ Ok(Transformed::no(plan))
+ }
+ }
+ // A correlated filter below these nodes cannot move above them:
+ // - A Window function reads the rows of its partition, so the
+ // filter changes the rows each function sees.
+ // - `DISTINCT ON` keeps the first row of each group, so the filter
+ // changes which row is first.
+ // - The pull up does not carry the columns of the pulled up filter
+ // through an Unnest, and fails with a schema error.
+ // - A RecursiveQuery repeats its recursive term.
+ // - The other nodes are not queries and are not in a subquery.
+ LogicalPlan::Window(_)
+ | LogicalPlan::Distinct(Distinct::On(_))
+ | LogicalPlan::Unnest(_)
Review Comment:
As the comment says, this stop works around a missing schema update, not a
semantic problem. An `Unnest` maps each input row to its own output rows, so a
filter on a column it carries keeps the same rows above or below it. Spark
allows `Generate` on a correlation path, and DuckDB pushes correlation through
`LOGICAL_UNNEST` the same way as through a filter.
The schema error comes from the `Unnest` keeping the schema it was built
with after `f_up` adds a column to the projection below it.
`LogicalPlan::recompute_schema` already rebuilds an `Unnest` through
`unnest_with_options`, so recomputing it in `f_up` might decorrelate the new
`subquery.slt` case instead of erroring. I have not tried it.
The stop also catches queries that select the correlated column next to
`unnest(...)`, for example `LATERAL (SELECT i.id, unnest(i.arr) AS u FROM i
WHERE i.id = o.k)`. The planner already carries `i.id` through the `Unnest`
there, so if I read it right `main` decorrelates it correctly and this PR
rejects 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]