Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ package org.apache.spark.sql.execution.exchange
import scala.collection.mutable
import scala.collection.mutable.ArrayBuffer

import org.apache.spark.internal.{LogKeys}
import org.apache.spark.internal.LogKeys
import org.apache.spark.sql.catalyst.expressions._
import org.apache.spark.sql.catalyst.plans._
import org.apache.spark.sql.catalyst.plans.physical._
Expand All @@ -30,7 +30,7 @@ import org.apache.spark.sql.connector.catalog.functions.Reducer
import org.apache.spark.sql.errors.QueryExecutionErrors
import org.apache.spark.sql.execution._
import org.apache.spark.sql.execution.datasources.v2.GroupPartitionsExec
import org.apache.spark.sql.execution.joins.{ShuffledHashJoinExec, SortMergeJoinExec}
import org.apache.spark.sql.execution.joins.{ShuffledHashJoinExec, ShuffledJoin, SortMergeJoinExec}
import org.apache.spark.sql.internal.SQLConf

/**
Expand Down Expand Up @@ -61,6 +61,10 @@ case class EnsureRequirements(
shuffleOrigin: ShuffleOrigin): Seq[SparkPlan] = {
assert(requiredChildDistributions.length == originalChildren.length)
assert(requiredChildOrderings.length == originalChildren.length)
// A storage-partitioned join handles its co-partitioning (and the projected join-key
// GroupPartitionsExec) separately in checkKeyGroupCompatible, so the projected-key grouping
// below must not run for it. For a non-join operator, do it inline here.
val isJoin = parent.exists(_.isInstanceOf[ShuffledJoin])

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This guard states a different condition from the one the comment appeals to, and the two sets come apart in both directions.

The comment says joins are excluded because checkKeyGroupCompatible handles their projection. But that helper matches only SortMergeJoinExec and ShuffledHashJoinExec, while the block that calls it is entered for any operator with two clustered children (parent.isDefined && children.length == 2 && childrenIndexes.length == 2).

One direction is currently harmless: SortMergeAsOfJoinExec is also a ShuffledJoin, so it is excluded here yet unhandled there. It still gets correct results, but by falling through to a plain shuffle rather than by anything this comment describes.

The other direction concerns me, and I could not finish confirming it - hence a question. CoGroupExec requires ClusteredDistribution on both children (objects.scala:638) and is not a ShuffledJoin, so this branch does run for it. children is then reassigned with the wrapped child, and specs a few lines down are computed from children(i).outputPartitioning - the projected partitioning. checkKeyGroupCompatible returns None, so areChildrenCompatible is false and each clustered child reaches withJoinKeyPositions(child, joinKeyPositions), which rewrites the node in place via g.copy(joinKeyPositions = Some(positions)). Those positions index the projected expression list, but GroupPartitionsExec applies them to its child's unprojected partitioning. For tables partitioned by (name, id) cogrouped on id: this branch computes [1] and inserts a node projecting to id; the block recomputes [0] against the now one-element list and overwrites, so the node projects position 0 of (name, id) and coalesces by name.

I did not verify the observable result. Reaching the overwrite also needs v2BucketingShuffleEnabled on, since it gates KeyedShuffleSpec.canCreatePartitioning - with the default false, bestSpecOpt is empty and the fallback shuffles and unwraps the node instead. Could you confirm whether a cogroup over two storage-partitioned tables with both configs on and a non-leading grouping key produces wrong groups?

Either way I'd gate on the structural fact rather than the trait: run this branch only when the operator will not enter that block, i.e. childrenIndexes.length <= 1. That is what the comment is really appealing to, it covers cogroups and as-of joins without enumerating join classes, and it keeps one owner of the projection per operator so a future ShuffledJoin or a new parent case in checkKeyGroupCompatible cannot move the boundary silently.

// Ensure that the operator's children satisfy their output distribution requirements.
var children = originalChildren.zip(requiredChildDistributions).map {
case (child, distribution) =>
Expand Down Expand Up @@ -107,8 +111,24 @@ case class EnsureRequirements(
}

case _ if groupedSatisfies.isDefined =>
// Grouped KeyedPartitioning already satisfies
child
// A grouped KeyedPartitioning already satisfies. However, when the operation keys
// are a strict subset of the partition keys (enabled via
// v2BucketingAllowKeysSubsetOfPartitionKeys), the partitions are still grouped by
// the full partition keys rather than by the operation keys, so a
// GroupPartitionsExec that projects to the operation keys must be inserted to
// coalesce partitions sharing the same operation key. `createShuffleSpec` computes
// exactly those projected positions when the config is enabled.
val kp = groupedSatisfies.get
distribution match {
case c: ClusteredDistribution if !isJoin =>
val spec = kp.createShuffleSpec(c).asInstanceOf[KeyedShuffleSpec]

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Only spec.joinKeyPositions is used, but with the config enabled createShuffleSpec also runs projectKeys(joinKeyPositions)._2 across every partition key and then .distinct on the result, to build a projectedPartitioning that this call site drops - and GroupPartitionsExec recomputes the same projection later from child.outputPartitioning. That is two O(number of partitions) passes per qualifying operator at planning time, on a path that exists precisely for tables partitioned finely enough to need coalescing.

KeyedShuffleSpec(kp, c).keyPositions gives the same information without the discarded projection, and reads more directly as "which partition expressions does this operation actually key on" than routing through a shuffle-spec factory. The cast also goes away with it.

spec.joinKeyPositions match {
case Some(positions) if positions != kp.expressions.indices.toSeq =>
GroupPartitionsExec(child, joinKeyPositions = Some(positions))
case _ => child
}
case _ => child
}

case _ if nonGroupedSatisfiesAsIs =>
// Non-grouped KeyedPartitioning satisfies without grouping
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4207,6 +4207,111 @@ class KeyGroupedPartitioningSuite extends DistributionAndOrderingSuiteBase with
}
}

test("window top-k over PARTITION BY subset of partition keys coalesces partitions") {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Consider adding an aggregate variant next to these. HashAggregateExec requires ClusteredDistribution(groupingExpressions), so SELECT id, sum(price) FROM items GROUP BY id on an (id, name)-partitioned table reaches this exact branch, and without the fix it emits one row per (id, name) partition instead of per id - a duplicated group, which is arguably a more familiar symptom than a ranking artifact. All three new tests go through ROW_NUMBER() OVER (PARTITION BY ...), so a reader could reasonably think the fix is window-specific.

A cogroup test would be worth more still, because that case is one the change newly affects rather than fixes: df1.groupByKey(...).cogroup(df2.groupByKey(...)) over two storage-partitioned tables is a two-clustered-child operator that is not a ShuffledJoin, so it takes the new branch and then also goes through the co-partitioning block. That is the interaction I asked about on the isJoin line.

// items is partitioned by (id, name). A top-k window that ranks by PARTITION BY id (a subset of
// the partition keys) must coalesce the (1,'aa') and (1,'bb') partitions before ranking so that
// id=1 is ranked across both rows and yields a single row; otherwise each partition is ranked
// independently and id=1 surfaces twice.
val items_partitions = Array(identity("id"), identity("name"))
createTable(items, itemsColumns, items_partitions)
sql(s"INSERT INTO testcat.ns.$items VALUES " +
s"(1, 'aa', 10.0, cast('2020-01-01' as timestamp)), " +
s"(1, 'bb', 20.0, cast('2020-01-01' as timestamp)), " +
s"(2, 'cc', 30.0, cast('2020-01-01' as timestamp))")

val query =
s"""SELECT id, name, price FROM (
| SELECT id, name, price, ROW_NUMBER() OVER (PARTITION BY id ORDER BY price DESC) rn
| FROM testcat.ns.$items
|) t WHERE rn = 1
|""".stripMargin
val expected = Seq(Row(1L, "bb", 20.0f), Row(2L, "cc", 30.0f))

// Result correctness does not depend on AQE: EnsureRequirements also runs in AQE's
// queryStagePreparationRules and likewise skips the GroupPartitionsExec, ranking id=1
// per-partition. Verify the wrong result under the default (AQE on) configuration.
withSQLConf(SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "true") {
checkAnswer(sql(query), expected)
}

// The plan-shape assertion needs a static, fully-planned tree, so disable AQE here.
withSQLConf(
SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "true") {
assert(collectAllGroupPartitions(sql(query).queryExecution.executedPlan).nonEmpty,
"GroupPartitionsExec expected to coalesce partitions sharing the subset key [id]")
}
}

test("window top-k over duplicated PARTITION BY key coalesces partitions") {
// Like the subset test above, but the window partition spec repeats the same key
// (PARTITION BY id, id). The projected positions that createShuffleSpec computes are still just
// [id], so the duplicated spec must not be mistaken for "no projection" and must still coalesce
// the (1,'aa') and (1,'bb') partitions.
val items_partitions = Array(identity("id"), identity("name"))
createTable(items, itemsColumns, items_partitions)
sql(s"INSERT INTO testcat.ns.$items VALUES " +
s"(1, 'aa', 10.0, cast('2020-01-01' as timestamp)), " +
s"(1, 'bb', 20.0, cast('2020-01-01' as timestamp)), " +
s"(2, 'cc', 30.0, cast('2020-01-01' as timestamp))")

val query =
s"""SELECT id, name, price FROM (
| SELECT id, name, price, ROW_NUMBER() OVER (PARTITION BY id, id ORDER BY price DESC) rn
| FROM testcat.ns.$items
|) t WHERE rn = 1
|""".stripMargin
val expected = Seq(Row(1L, "bb", 20.0f), Row(2L, "cc", 30.0f))

withSQLConf(SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "true") {
checkAnswer(sql(query), expected)
}

withSQLConf(
SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "true") {
assert(collectAllGroupPartitions(sql(query).queryExecution.executedPlan).nonEmpty,
"GroupPartitionsExec expected to coalesce partitions sharing the duplicated key [id]")
}
}

test("window top-k over union output partitioning coalesces partitions") {
// t1 and t2 are both partitioned by (id, name). With union output partitioning enabled, the
// union reports a KeyedPartitioning over (id, name), so a top-k window over PARTITION BY id (a
// strict subset) must still coalesce the (1,'aa') and (1,'bb') partitions coming from t1.
val partitions = Array(identity("id"), identity("name"))
withTable("t1", "t2") {
createTable("t1", itemsColumns, partitions)
sql("INSERT INTO testcat.ns.t1 VALUES " +
"(1, 'aa', 10.0, cast('2020-01-01' as timestamp)), " +
"(1, 'bb', 20.0, cast('2020-01-01' as timestamp))")
createTable("t2", itemsColumns, partitions)
sql("INSERT INTO testcat.ns.t2 VALUES (2, 'cc', 30.0, cast('2020-01-01' as timestamp))")

val query =
"""SELECT id, name, price FROM (
| SELECT id, name, price,
| ROW_NUMBER() OVER (PARTITION BY id ORDER BY price DESC) rn
| FROM (
| SELECT id, name, price FROM testcat.ns.t1
| UNION ALL
| SELECT id, name, price FROM testcat.ns.t2
| )
|) t WHERE rn = 1
|""".stripMargin
val expected = Seq(Row(1L, "bb", 20.0f), Row(2L, "cc", 30.0f))

withSQLConf(
SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "true",
SQLConf.UNION_OUTPUT_PARTITIONING.key -> "true") {
checkAnswer(sql(query), expected)
assert(collectAllGroupPartitions(sql(query).queryExecution.executedPlan).nonEmpty,
"GroupPartitionsExec expected to coalesce union partitions sharing the key [id]")
}
}
}

test("SPARK-57881: storage-partitioned join leverages union output KeyedPartitioning to " +
"avoid shuffle") {
val cols = Array(
Expand Down