ulysses-you commented on code in PR #57603: URL: https://github.com/apache/spark/pull/57603#discussion_r3670708378
########## sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/PullUpProjectAliasThroughWindow.scala: ########## @@ -0,0 +1,122 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.spark.sql.catalyst.optimizer + +import org.apache.spark.sql.catalyst.expressions.{Alias, Attribute, AttributeMap, AttributeSet, NamedExpression} +import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan, Project, Window} +import org.apache.spark.sql.catalyst.rules.Rule +import org.apache.spark.sql.catalyst.trees.TreePattern.WINDOW + +/** + * A [[Window]] operator does not project its output partitioning or ordering through aliases: its + * physical `outputPartitioning`/`outputOrdering` are pure pass-throughs of the child's. So when a + * window partitioned/ordered by `k` is followed by a consumer (aggregate, repartition, a chained + * window, ...) that requires distribution/ordering on a rename `k AS a`, a redundant shuffle and/or + * sort is inserted: the window already shuffled/sorted on `k`, but `HashPartitioning(k)` does not + * satisfy `ClusteredDistribution(a)` because the rename lives in a [[Project]] *below* the window + * while the analyzer-inserted [[Project]] *above* references `a` as a bare [[Attribute]] (empty + * alias map), and `k` and `a`, though value-identical, have distinct expr ids. + * + * This rule pulls such aliases up from the bottom [[Project]] into the top [[Project]], across the + * chain of one or more [[Window]] operators in between (`Project - Window+ - Project`; adjacent + * windows that cannot be collapsed, e.g. with different order specs, leave no [[Project]] between + * them). The top project's alias map then maps `k -> a`, and the + * `PartitioningPreservingUnaryExecNode` / `OrderPreservingUnaryExecNode` machinery that + * `ProjectExec` mixes in projects the windows' `HashPartitioning(k)`/`SortOrder(k)` up through the + * alias, satisfying the consumer. The [[Window]] operators themselves are left untouched. + * + * Besides removing the redundant shuffle/sort, this also narrows the data crossing the window's + * shuffle: the alias no longer flows through the window as a separate column (only its input `k`, + * already needed for partitioning, does) and is recomputed cheaply in the top project above the + * exchange. + * + * An entry of the bottom project is pulled up only when: + * - it is an [[Alias]] (bare pass-through attributes stay below so the windows keep producing + * them for the top project to reference); + * - no window in the chain references it, so window semantics are unaffected; + * - all of its input attributes remain produced by the pruned bottom project (i.e. they are + * referenced by some window and thus retained), so a computed alias whose inputs would be + * dropped is not lifted above windows that no longer output them; + * - the top project consumes its output solely as a bare pass-through attribute (never inside a + * larger expression), so replacing that attribute with the alias fully preserves it. + * + * The rewrite is a no-op on its own output (a moved alias is no longer a bare attribute above nor + * present below), so it converges immediately at the fixed point. + */ +object PullUpProjectAliasThroughWindow extends Rule[LogicalPlan] { + + /** + * Matches a non-empty chain of [[Window]] operators bottomed out by a [[Project]], returning the + * windows top-to-bottom and that bottom project. + */ + private object WindowChain { + def unapply(plan: LogicalPlan): Option[(Seq[Window], Project)] = plan match { + case w @ Window(_, _, _, lower: Project, _) => Some((Seq(w), lower)) + case w @ Window(_, _, _, WindowChain(windows, lower), _) => Some((w +: windows, lower)) + case _ => None + } + } + + override def apply(plan: LogicalPlan): LogicalPlan = plan.transformWithPruning( + _.containsPattern(WINDOW), ruleId) { + // Match the `Project - Window+ - Project` shape: the top project is the windows' parent + // scaffolding, and the bottom project is where the rename (`key AS userid`) is defined. + case p @ Project(projectList, WindowChain(windows, lower)) => + // Entries any window references must stay below: they define what the bottom project must + // keep producing, and in particular carry the partition/order key attributes the pulled-up + // aliases reuse. + val windowRefs = AttributeSet(windows.flatMap(_.references)) + val (retained, candidates) = + lower.projectList.partition(e => windowRefs.contains(e.toAttribute)) + val retainedAttrs = AttributeSet(retained.map(_.toAttribute)) + // Attributes the top project consumes as bare pass-through entries, and those it consumes + // inside a larger expression (an alias child, a window-output reference, ...). + val bareAttrs = AttributeSet(projectList.collect { case a: Attribute => a }) + val referencedByExpr = AttributeSet(projectList.flatMap { + case _: Attribute => Nil + case other => other.references + }) + // Pull up an alias only when: its inputs survive in the pruned bottom project, the top + // project passes its output through as a bare attribute, and the top project does not also + // consume that output inside an expression (which would require the windows to keep it). + val pullUp = candidates.collect { + case a: Alias + if a.references.subsetOf(retainedAttrs) && Review Comment: Good catch, fixed in e16b666. Added an `a.deterministic` guard to the pull-up condition so nondeterministic aliases stay below the window, plus a "no rewrite for a nondeterministic alias" regression test using `spark_partition_id()`. ########## sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/PullUpProjectAliasThroughWindow.scala: ########## @@ -0,0 +1,122 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.spark.sql.catalyst.optimizer + +import org.apache.spark.sql.catalyst.expressions.{Alias, Attribute, AttributeMap, AttributeSet, NamedExpression} +import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan, Project, Window} +import org.apache.spark.sql.catalyst.rules.Rule +import org.apache.spark.sql.catalyst.trees.TreePattern.WINDOW + +/** + * A [[Window]] operator does not project its output partitioning or ordering through aliases: its + * physical `outputPartitioning`/`outputOrdering` are pure pass-throughs of the child's. So when a + * window partitioned/ordered by `k` is followed by a consumer (aggregate, repartition, a chained + * window, ...) that requires distribution/ordering on a rename `k AS a`, a redundant shuffle and/or + * sort is inserted: the window already shuffled/sorted on `k`, but `HashPartitioning(k)` does not + * satisfy `ClusteredDistribution(a)` because the rename lives in a [[Project]] *below* the window + * while the analyzer-inserted [[Project]] *above* references `a` as a bare [[Attribute]] (empty + * alias map), and `k` and `a`, though value-identical, have distinct expr ids. + * + * This rule pulls such aliases up from the bottom [[Project]] into the top [[Project]], across the + * chain of one or more [[Window]] operators in between (`Project - Window+ - Project`; adjacent + * windows that cannot be collapsed, e.g. with different order specs, leave no [[Project]] between + * them). The top project's alias map then maps `k -> a`, and the + * `PartitioningPreservingUnaryExecNode` / `OrderPreservingUnaryExecNode` machinery that + * `ProjectExec` mixes in projects the windows' `HashPartitioning(k)`/`SortOrder(k)` up through the + * alias, satisfying the consumer. The [[Window]] operators themselves are left untouched. + * + * Besides removing the redundant shuffle/sort, this also narrows the data crossing the window's + * shuffle: the alias no longer flows through the window as a separate column (only its input `k`, + * already needed for partitioning, does) and is recomputed cheaply in the top project above the + * exchange. + * + * An entry of the bottom project is pulled up only when: + * - it is an [[Alias]] (bare pass-through attributes stay below so the windows keep producing + * them for the top project to reference); + * - no window in the chain references it, so window semantics are unaffected; + * - all of its input attributes remain produced by the pruned bottom project (i.e. they are + * referenced by some window and thus retained), so a computed alias whose inputs would be + * dropped is not lifted above windows that no longer output them; + * - the top project consumes its output solely as a bare pass-through attribute (never inside a + * larger expression), so replacing that attribute with the alias fully preserves it. + * + * The rewrite is a no-op on its own output (a moved alias is no longer a bare attribute above nor + * present below), so it converges immediately at the fixed point. + */ +object PullUpProjectAliasThroughWindow extends Rule[LogicalPlan] { + + /** + * Matches a non-empty chain of [[Window]] operators bottomed out by a [[Project]], returning the + * windows top-to-bottom and that bottom project. + */ + private object WindowChain { + def unapply(plan: LogicalPlan): Option[(Seq[Window], Project)] = plan match { + case w @ Window(_, _, _, lower: Project, _) => Some((Seq(w), lower)) + case w @ Window(_, _, _, WindowChain(windows, lower), _) => Some((w +: windows, lower)) + case _ => None + } + } + + override def apply(plan: LogicalPlan): LogicalPlan = plan.transformWithPruning( + _.containsPattern(WINDOW), ruleId) { + // Match the `Project - Window+ - Project` shape: the top project is the windows' parent + // scaffolding, and the bottom project is where the rename (`key AS userid`) is defined. + case p @ Project(projectList, WindowChain(windows, lower)) => + // Entries any window references must stay below: they define what the bottom project must + // keep producing, and in particular carry the partition/order key attributes the pulled-up + // aliases reuse. + val windowRefs = AttributeSet(windows.flatMap(_.references)) + val (retained, candidates) = + lower.projectList.partition(e => windowRefs.contains(e.toAttribute)) + val retainedAttrs = AttributeSet(retained.map(_.toAttribute)) + // Attributes the top project consumes as bare pass-through entries, and those it consumes + // inside a larger expression (an alias child, a window-output reference, ...). + val bareAttrs = AttributeSet(projectList.collect { case a: Attribute => a }) + val referencedByExpr = AttributeSet(projectList.flatMap { + case _: Attribute => Nil + case other => other.references + }) + // Pull up an alias only when: its inputs survive in the pruned bottom project, the top + // project passes its output through as a bare attribute, and the top project does not also + // consume that output inside an expression (which would require the windows to keep it). + val pullUp = candidates.collect { + case a: Alias + if a.references.subsetOf(retainedAttrs) && + bareAttrs.contains(a.toAttribute) && + !referencedByExpr.contains(a.toAttribute) => a + } + if (pullUp.isEmpty) { + p + } else { + val pullUpSet: Set[NamedExpression] = pullUp.toSet + // Key the lookup by expr id (not by `Attribute.equals`, which also compares qualifier): + // the top project may reference these attributes with a different qualifier than the alias + // carries below (e.g. across a subquery alias). + val pullUpMap = AttributeMap(pullUp.map(a => a.toAttribute -> a)) + val newProjectList = projectList.map { + case attr: Attribute => pullUpMap.getOrElse(attr, attr) Review Comment: Thanks, fixed in e16b666. The lifted alias is now rebuilt with the top attribute's name via `Alias.withName(attr.name)`, which preserves the child, expr id, qualifier, and metadata while keeping the top attribute's cosmetic name, so the output schema is unchanged. Added a regression test that references the top attribute under a different (expr-id-equivalent) name and asserts the output name is preserved. -- 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]
