From f95e44a9f422e7e251f03173f25040c60ca698ac Mon Sep 17 00:00:00 2001 From: Xiduo You Date: Sun, 20 Sep 2026 15:15:46 +0800 Subject: [PATCH 1/2] [SPARK-59671][SQL] Validate a co-partitioned pair by its children's pairing ### What changes were proposed in this pull request? `ValidateRequirements` asks every child of a clustered operator to satisfy its distribution on its own. A `ClusteredDistribution` is the one distribution an operator can owe two children together rather than one by one, so an operator whose children all owe one is judged on their mutual layout, and a pair aligned without grouping (which partially clustered distribution builds on purpose) is one the sides agree on while neither is grouped. - The children of such an operator are judged together, on the layouts they report (`PartitioningCollection.specsForPairing`, through its `reportedSpecOf`), and one member of the first side has to pair with every other side. That is every multi-child clustered operator, not only a join: one that zips corresponding partitions, a cogroup for instance, reads a layout both children have to hold together. Every other child, an operator with a single clustered child included, keeps the per-side check. - What a side offers is the layout it reports, restricted to the partition expressions that carry a cluster key: no key is deduped and none is re-sorted, so the count and the order are the member's own. `KeyedPartitioning.createShuffleSpec` is the layout a node would emit and is deliberately not what the validator reads: it dedups and sorts, and its count is the grouped one, which the plan does not hold. - The one member a finished plan may report without satisfying the distribution is the ungrouped shape partially clustered distribution spreads, and one producer builds it: `EnsureRequirements.checkKeyGroupCompatible`, which is entered for a sort-merge or shuffled-hash join. The waiver is confined to that configuration and that producer, and the answer is passed down to the pairing so it cannot be read one way in one place and the other way in the next. What is left to the member itself, its partition count and the permission for the collapse it went through, is asked in the spec: a member whose keys were collapsed is admitted ungrouped only where the planner would agree to group it, and an ungrouped member keeps the keys it reports. - What the pairing cannot say is how the two sides hold a key's rows: a spread side and one that repeats the whole group report the same keys as two sides that split the key between them, and no layout distinguishes those, so that rests on the producer, which is the join path. Against the planner's own admission of a member (`EnsureRequirements.createKeyedShuffleSpecs`), the coverage of every operation key (`spark.sql.requireAllClusterKeysForCoPartition`) is not asked here: it is a skew heuristic, and a member whose partitioning keys cover only a subset of the operation's keys is a sound pairing. ### Why are the changes needed? A storage-partitioned join planned by partially clustered distribution aligns its sides without grouping either of them: the side that keeps its splits spreads them, the other replicates its group across them, so both report keys that repeat on purpose. Such a pair fails the per-side check at the join node even though the two sides agree key by key, and `AdaptiveSparkPlanExec.optimizeQueryStage` validates a stage's whole candidate plan before accepting an `AQEShuffleReadRule` change, so every shuffle read in that stage stays uncoalesced, unrelated ones included. Measured on `d39cc1784c0` (base) versus this head, both with `spark.sql.sources.v2.bucketing.partiallyClusteredDistribution.enabled=true`: | Probe | base | head | |---|---|---| | A pair the rule planned (`EnsureRequirements.apply`) | `ValidateRequirements.validate` false | true | | The aggregate's shuffle read in the join's stage | not coalesced (`AQEShuffleReadExec.hasCoalescedPartition` false) | coalesced | | A three-table chain whose outer join reads a projection over both key columns | `validate` false | true, shuffle-free, still coalesced | The planner-side half of this hazard landed in SPARK-59272, which declines a pairing whose sides no longer declare the same aligned key sequence: that closes the pairs a `GroupPartitionsExec` gives up on, which should not be built at all. A pair aligned without grouping is the other half, and it is built on purpose, so the validator has to read the pairing instead of each child on its own. SPARK-59688, merged in the meantime, tightened `KeyedShuffleSpec.isCompatibleWith` to answer for a pair as it stands, which is the question asked here; a pair that took the reduce still answers it through `hasSameReducedKeys`, so the shapes in this PR are unaffected. ### Does this PR introduce _any_ user-facing change? No. No plan changes unless the plan already contains such an alignment, and no new configuration. ### How was this patch tested? New tests, by what each one is there for. Those that assert the new acceptance pass only with the main-code change, since the base asked `satisfies` of each child on its own and an ungrouped keyed child fails that in every configuration; those that assert a refusal hold on the base as well, which is where the guards are. - `ValidateRequirementsSuite` - the exemption and its refusals: a keyed pair that repeats its keys position by position passes where partially clustered distribution is on, while key sets that disagree, a differing key order, a hashed side, a lone clustered child and two sides that each satisfy on their own but do not line up all fail; - the rule's own pair: a pair planned by `EnsureRequirements` under partially clustered distribution passes; - the producer boundary: the same ungrouped pair a shuffled-hash join reads is refused for an as-of join, which is a `ShuffledJoin` too and builds no spread side; - the collection shape: a side reporting several keyed alternatives is judged on whichever of them pairs, not on the first one, and a side offering no member keyed on the join keys does not pair; - a multi-child clustered operator that is not a join: a cogroup over the pair a join reads is refused, and the mutual check on the layouts its sides report refuses a differing partition count and a differing key order while accepting two sides that agree; - the key a side is judged on: a side partitioned on `[a, b]` serving an operation on `[a]` is judged on the key the operation clusters on, as it reports it, so the pair stands and a side whose keys run the other way is refused; - the admission's edges: an ungrouped side is refused where nothing builds one; a pair is judged on the layouts its sides report rather than on a projection no node made, for a key order and for a partition count; a collapsed pair is refused while its collapse may not be grouped, partial clustering on or not; a lone clustered child still owes its own grouping; and a pair that lines up still owes its operator an ordering. - `KeyGroupedPartitioningSuite` - with partially clustered distribution on, the aggregate's shuffle read in the stage holding a two-table join coalesces, which the base blocks by refusing the pair; the test also pins that the join side shuffles nothing, that the chain shuffles once, and that the join shares the final stage with that read, which is what makes the coalesce a decision the join can block; - a three-table chain whose outer join reads a projection that keeps both key columns stays shuffle-free, passes validation and keeps its coalescing. - `ShuffleSpecSuite`: the specs a side offers are the layouts it reports: a keyed member offers its own partitions under the operation's key, keeping its count and the order it reports; an ungrouped one is offered only where something builds one and produces it; a member whose partitioning keys are a superset of the operation's keys offers its own keys rather than a projection onto them; a marked member is not offered through a projection; a member whose keys do not cover the clustering offers nothing; a member that is not keyed offers its own spec; a collapsed member the planner would not group is not offered ungrouped; and a count the operation pinned is asked as it stands. Ran locally, all green: `ShuffleSpecSuite` and `DistributionSuite` (52 tests), and in one run `ValidateRequirementsSuite`, `EnsureRequirementsSuite`, `GroupPartitionsExecSuite`, `KeyGroupedPartitioningSuite` and `ProjectedOrderingAndPartitioningSuite` (362 tests), with `dev/lint-scala` clean for catalyst and sql, main and test sources. ### Was this patch authored or co-authored using generative AI tooling? Assisted-by: Qwen 3.8 Flash --- .../plans/physical/partitioning.scala | 129 ++++++- .../spark/sql/catalyst/ShuffleSpecSuite.scala | 74 +++- .../exchange/EnsureRequirements.scala | 18 +- .../exchange/ValidateRequirements.scala | 78 +++- .../sql/execution/joins/ShuffledJoin.scala | 17 + .../KeyGroupedPartitioningSuite.scala | 98 ++++- .../exchange/ValidateRequirementsSuite.scala | 345 +++++++++++++++++- 7 files changed, 716 insertions(+), 43 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/physical/partitioning.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/physical/partitioning.scala index 187163cfc0b6..c08e4ee87d31 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/physical/partitioning.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/physical/partitioning.scala @@ -564,8 +564,8 @@ case class KeyLayout( * * == Distribution Satisfaction and Grouping == * Besides the default `satisfies()`, `KeyedPartitioning` answers a family of questions. They differ - * in what they let happen to the data before the distribution counts as met. Only - * `keysMaySatisfy()` is asked from outside the class. The rest build it and `satisfies()` up. + * in what they let happen to the data before the distribution counts as met. `keysMaySatisfy()` and + * `mayServeUngrouped()` are asked from outside the class. The rest build them and `satisfies()` up. * * - `keysSatisfy()`: do the keys as they stand co-locate every cluster key, with nothing left for * a node to project away? This is the strict question, and `satisfies()` is it plus `isGrouped`. @@ -578,6 +578,9 @@ case class KeyLayout( * distribution? It is asked of non-grouped partitionings only, where such a node coalesces the * duplicate keys. So the answer is `keysCanSatisfy()`, plus whether that coalescing is allowed. * "Key Collapse" below says when it is not. + * - `mayServeUngrouped()`: the same question for a partitioning a finished plan reports ungrouped, + * where nobody is left to insert that node. So the answer is `keysSatisfy()`, plus the same + * coalescing permission, and the projection half of `keysCanSatisfy()` is out. * - `keysMaySatisfy()`: the same question for a partitioning that may already be grouped, which is * what `EnsureRequirements` asks. It is `mayGroupToSatisfy()` for a non-grouped one, and * `keysCanSatisfy()` for a grouped one, which has no duplicate keys left for the node to @@ -937,6 +940,10 @@ case class KeyedPartitioning( * side whose keys run the wrong way still answers `true` here. `EnsureRequirements.resolveChild` * is what compares the keys against the required ordering and builds the sorting node, and it has * to stay: nothing below tells it the keys are already in order. + * + * `mayServeUngrouped` is the caller that asks this outside the class: a member a finished plan + * reports ungrouped has nothing left to group it, so what it holds is what the question is about. + * See there for the permission it adds on top. */ private def keysSatisfy(required: Distribution): Boolean = { required match { @@ -990,14 +997,30 @@ case class KeyedPartitioning( satisfiesAfterProjection || keysSatisfy(required) } + /** + * The `Key Collapse` permission: whether a collapsed layout may be grouped at all. Read wherever + * a grouping node or its absence has to be judged. See the class doc's `Key Collapse` section for + * what the flag records. + */ + private def collapsedLayoutMayBeGrouped: Boolean = + !isCollapsed || SQLConf.get.v2BucketingAllowKeysSubsetOfPartitionKeys + + /** + * Whether this partitioning, reported ungrouped, may still serve `required` as it stands: the + * keys it holds co-locate every cluster key, and the collapse a grouping node would settle is + * permitted. `mayGroupToSatisfy` is the planner's question, which also admits a member such a + * node would serve; a finished plan has nobody left to insert one, and this is what it leaves. + */ + private[sql] def mayServeUngrouped(required: ClusteredDistribution): Boolean = + collapsedLayoutMayBeGrouped && keysSatisfy(required) + /** * Ask this only of a partitioning that is not grouped, since a grouped one has nothing to * coalesce and `keysMaySatisfy` covers both. See the class doc. */ private[sql] def mayGroupToSatisfy(required: Distribution): Boolean = { val mayCoalesce = required match { - case _: ClusteredDistribution => - !isCollapsed || SQLConf.get.v2BucketingAllowKeysSubsetOfPartitionKeys + case _: ClusteredDistribution => collapsedLayoutMayBeGrouped case _ => true } // The permission is the cheap half, so it is asked first. @@ -1371,8 +1394,9 @@ case class PartitioningCollection(partitionings: Seq[Partitioning]) // operation's still can, through the projection a `GroupPartitionsExec` performs. That is the // admission set this filter had before `satisfies` became strict, up to one shape it now also // keeps, a partition expression that *is* a cluster key, which `areKeysCompatible` turns away - // anyway. The set matters because `ValidateRequirements` builds a spec from a finished plan - // through here. + // anyway. The set matters because the planner's `shuffleToCoPartition` reads it through here to + // pick the layout the other children are laid out on, and a member this filter drops is one the + // collection cannot offer as that layout. // // Every admitted member stays, because `isCompatibleWith` answers for any of them and the // collection cannot know which one the other side matched. The cost is that @@ -1420,10 +1444,13 @@ object PartitioningCollection { * stands, and only a keyed partitioning answers the two differently. * * A partitioning that is not grouped is not admitted, even though a node would also group it. - * This is the admission set `satisfies` gave the one caller before it became strict, and the - * caller feeds `ValidateRequirements` as well as the planner, so it does not widen what a - * finished plan is checked against. `EnsureRequirements.createKeyedShuffleSpecs` is where the - * planner asks the wider question, for a child it is about to group itself. + * This is the admission set `satisfies` gave the caller before it became strict, and it is what + * `PartitioningCollection.createShuffleSpec` answers with, which the planner's + * `shuffleToCoPartition` reads to lay the other children out on what this one reports. + * `ValidateRequirements` does not read it: `specsForPairing` is its admission, and it admits an + * ungrouped member only where partially clustered distribution reports one. + * `EnsureRequirements.createKeyedShuffleSpecs` is where the planner asks the wider question, for + * a child it is about to group itself. * * The partition count clause is here for consistency with the strict question. A node changes the * count, so the pre-grouping one is no prediction of it. @@ -1437,6 +1464,74 @@ object PartitioningCollection { case other => other.satisfies(required) } + /** + * The specs `p` offers for `distribution`, one per member that may serve it, each of them the + * layout that member reports (`reportedSpecOf`): the partitions it has, in the order it reports + * them, with the partition expressions that carry no cluster key left out. Leaving those out is + * what lets two members pair when the operation clusters on part of what they are partitioned + * on, and it invents nothing: a kept key is that partition's own, the count is the member's own, + * and the order is the member's own. + * + * A keyed member is admitted on `satisfies`, the as-it-stands question, count included; a member + * that is not keyed is asked on `satisfies` as well and has no projection to make. The one member + * whose layout a finished plan may report without satisfying the distribution is an ungrouped + * keyed one: that is the shape partially clustered distribution spreads. It is admitted only + * where the caller says so, through `mayBeUngrouped`, which carries both halves of that + * admission: the configuration that builds such a shape, and the operator whose producer spreads + * it. `mayServeUngrouped` is the as-it-stands question for it, since `satisfies` + * adds `isGrouped` on top. + * + * This is the planner's admission of a member (`EnsureRequirements.createKeyedShuffleSpecs`) less + * the coverage of every operation key it requires there + * (`spark.sql.requireAllClusterKeysForCoPartition`), which is a skew heuristic: a member whose + * partitioning keys cover only a subset of the operation's keys is a sound pairing. + * + * @param mayBeUngrouped whether an ungrouped keyed member of `p` may serve `distribution` here. + * Only the caller knows both halves of that; see above. + */ + private[sql] def specsForPairing( + p: Partitioning, + distribution: ClusteredDistribution, + mayBeUngrouped: Boolean): Seq[ShuffleSpec] = + flatten(p).flatMap { + case k: KeyedPartitioning => + val pairsAsIs = k.satisfies(distribution) || + (mayBeUngrouped && + distribution.requiredNumPartitions.forall(_ == k.numPartitions) && + k.mayServeUngrouped(distribution)) + Option.when(pairsAsIs)(reportedSpecOf(k, distribution)) + case other => + Option.when(other.satisfies(distribution))(other.createShuffleSpec(distribution)) + } + + /** + * The spec for a member a finished plan reports: its own layout with the partition expressions + * that carry no cluster key left out. No key is deduped and no key is re-sorted, so a side offers + * a view of the partitions the plan has rather than the layout + * `KeyedPartitioning.createShuffleSpec` builds for a node a planner would insert: that one dedups + * and sorts, and its count is the grouped one, which the plan does not hold. + * + * A projection here is a relabelling of the member's partitions, and this leans on the caller for + * that: only `keysSatisfy`'s projecting half admits such a member, and it holds the projection to + * one that merges no partition, handing the positions it checked to this method. + * + * A marked layout is not projected: its claim routes the rows of an undeclared key by a hash over + * these keys in this order, and dropping one breaks the claim. The clause cannot fire as the + * admission stands, and it is kept because the causality is not local: an admitted marked member + * has its keys covering the clustering structurally, since the projecting half of `keysSatisfy` + * excludes a marked one, so it leaves here unprojected either way. + */ + private def reportedSpecOf( + k: KeyedPartitioning, + distribution: ClusteredDistribution): KeyedShuffleSpec = { + val positions = k.positionsCoveringClusterKeys(distribution).toSeq + if (positions.size == k.expressions.length || k.mayContainUnknownPartitionKeys) { + KeyedShuffleSpec(k, distribution) + } else { + KeyedShuffleSpec(k.project(positions), distribution, Some(positions)) + } + } + /** * The uniform `mayContainUnknownPartitionKeys` marker of `p`'s keyed members, read from its * first keyed member, the same one `checkKeyedPartitioningInvariant` compares against. @@ -1639,12 +1734,14 @@ case object SinglePartitionShuffleSpec extends LeafShuffleSpec { // disagree when the subset config projects them onto different key sets. // // `EnsureRequirements` never reaches this: a spec whose `canCreatePartitioning` is false is - // never the best one. The one production caller that can put a collection on the `other` side - // is `ValidateRequirements`' `specs.tail.forall(_.isCompatibleWith(specs.head))`, and there the - // stricter answer is the safer one. It does leave this direction stricter than the collection's - // own `exists`, against the symmetry this trait's doc assumes, but nothing observable follows: - // only `KeyedShuffleSpec` can make members disagree on `numPartitions`, and it has no - // `SinglePartitionShuffleSpec` case, so that direction is already false. + // never the best one. Nor does any other production caller put a collection on the `other` + // side: `pickCoPartitionTarget` flattens before it pairs a member, and `ValidateRequirements` + // offers leaf specs and pairs those, where it used to compare every child against the first + // one's whole spec. So what is left is the answer `ShuffleSpecSuite` pins. It stays `forall` + // rather than the `exists` a collection answers with, against the symmetry the trait doc + // assumes, but nothing observable follows: only `KeyedShuffleSpec` can make members disagree on + // `numPartitions`, and it has no `SinglePartitionShuffleSpec` case, so that direction is + // already false. case ShuffleSpecCollection(specs) => specs.forall(isCompatibleWith) } diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/ShuffleSpecSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/ShuffleSpecSuite.scala index 86eec4be3870..cbefa271ad4d 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/ShuffleSpecSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/ShuffleSpecSuite.scala @@ -93,9 +93,8 @@ class ShuffleSpecSuite extends SparkFunSuite with SQLHelper { test("SPARK-59289: createShuffleSpec drops a keyed member that is not grouped") { val a = AttributeReference("a", IntegerType)() val clustered = ClusteredDistribution(Seq(a)) - // A node would group this one too, but the admission set here is what a finished plan is - // checked against by `ValidateRequirements`, so it stays what the strict question admitted - // before a projection was allowed to answer it. + // A node would group this one too, but a collection is a layout the plan holds, so it stays + // what the strict question admitted before a projection was allowed to answer it. val ungrouped = KeyedPartitioning(Seq(a), Seq(InternalRow(1), InternalRow(1), InternalRow(2))) val grouped = KeyedPartitioning(Seq(a), Seq(InternalRow(1), InternalRow(2), InternalRow(3))) assert(!ungrouped.isGrouped && grouped.isGrouped, "test setup") @@ -114,6 +113,75 @@ class ShuffleSpecSuite extends SparkFunSuite with SQLHelper { assert(specs.head.isInstanceOf[HashShuffleSpec], s"and it is the hash one, got $specs") } + test("SPARK-59671: the specs a side offers are the layouts it reports") { + val a = AttributeReference("a", IntegerType)() + val b = AttributeReference("b", IntegerType)() + val cd = ClusteredDistribution(Seq(a)) + // An ungrouped keyed member is the layout partially clustered distribution spreads, and it is + // offered as the layout it holds: nobody is left to group it. Only a caller that knows both + // halves of that admission, the configuration and the operator whose producer spreads one, says + // so; the helper takes the answer rather than reading a configuration for itself. + val spread = KeyedPartitioning(Seq(a), Seq(InternalRow(1), InternalRow(1), InternalRow(2))) + val specs = PartitioningCollection.specsForPairing(spread, cd, mayBeUngrouped = true) + assert(specs.size == 1, s"one member, one spec, got $specs") + val spec = specs.head.asInstanceOf[KeyedShuffleSpec] + assert((spec.partitioning eq spread) && spec.joinKeyPositions.isEmpty, + s"its own layout, no projection to make, got $specs") + + // Everywhere else a member that only a grouping node would serve is not offered at all. + assert(PartitioningCollection.specsForPairing(spread, cd, mayBeUngrouped = false).isEmpty, + "an ungrouped member is a plan only where something builds one and produces it") + + // The permission for a collapse is the planner's, and it survives into this admission: a + // collapsed member is not offered ungrouped where the planner would refuse to group it either. + val collapsed = spread.withLayout(_.copy(isCollapsed = true)) + withSQLConf(SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "false") { + assert(PartitioningCollection + .specsForPairing(collapsed, cd, mayBeUngrouped = true).isEmpty, + "a collapsed member the planner would not group is not offered ungrouped") + } + + // The subset permission applies where the operation's keys are a subset of the member's + // partitioning keys, so a member may carry an expression the operation does not cluster on. It + // is offered as its own partitions under the key the operation clusters on: that expression is + // left out, and no key is deduped or re-sorted for it. + withSQLConf(SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "true") { + val source = KeyedPartitioning(Seq(a, b), Seq(InternalRow(1, 1), InternalRow(2, 2))) + val offered = PartitioningCollection.specsForPairing(source, cd, mayBeUngrouped = false) + .head.asInstanceOf[KeyedShuffleSpec] + assert(offered.partitioning.expressions === Seq(a) && + offered.partitioning.numPartitions === source.numPartitions && + offered.joinKeyPositions === Some(Seq(0)), + s"the member's own partitions under the operation's key, got $offered") + + // A marked layout is never offered through a projection: the projecting half of `keysSatisfy` + // excludes a marked member, so the keys an offered marked one carries are the keys its claim + // is over. + val marked = source.withLayout(_.copy(mayContainUnknownPartitionKeys = true)) + assert(PartitioningCollection.specsForPairing(marked, cd, mayBeUngrouped = false).isEmpty, + "a marked member is not offered through a projection") + } + + // A member whose keys do not cover the clustering is not offered, and a member that is not + // keyed is asked for its own spec. + val elsewhere = KeyedPartitioning(Seq(b), Seq(InternalRow(1), InternalRow(2))) + assert(PartitioningCollection.specsForPairing(elsewhere, cd, mayBeUngrouped = false).isEmpty) + assert(PartitioningCollection + .specsForPairing(HashPartitioning(Seq(a), 2), cd, mayBeUngrouped = false).head + .isInstanceOf[HashShuffleSpec]) + + // A count the operation pinned is asked as it stands: a member of another size is not offered + // even where its keys cover the clustering. + val four = KeyedPartitioning(Seq(a), (1 to 4).map(InternalRow(_))) + def pairingFour(pinned: Int): Seq[ShuffleSpec] = + PartitioningCollection.specsForPairing( + four, ClusteredDistribution(Seq(a), requiredNumPartitions = Some(pinned)), + mayBeUngrouped = false) + assert(pairingFour(4).size == 1, + "the size the operation asks for is the one the member reports") + assert(pairingFour(3).isEmpty, "a member of another size does not serve the distribution") + } + protected def checkCompatible( left: ShuffleSpec, right: ShuffleSpec, diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/EnsureRequirements.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/EnsureRequirements.scala index 1eea8a006ed2..155a88c8eae0 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/EnsureRequirements.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/EnsureRequirements.scala @@ -591,13 +591,8 @@ case class EnsureRequirements( leftRequired: ClusteredDistribution, right: SparkPlan, rightRequired: ClusteredDistribution): Option[Seq[SparkPlan]] = { - parent match { - case smj: SortMergeJoinExec => - checkKeyGroupCompatible(left, leftRequired, right, rightRequired, smj.joinType) - case sj: ShuffledHashJoinExec => - checkKeyGroupCompatible(left, leftRequired, right, rightRequired, sj.joinType) - case _ => - None + ShuffledJoin.partiallyClusteredJoinType(parent).flatMap { joinType => + checkKeyGroupCompatible(left, leftRequired, right, rightRequired, joinType) } } @@ -904,6 +899,15 @@ case class EnsureRequirements( } // Now we need to push-down the common partition information to the `GroupPartitionsExec`s. + // + // The two arguments below say which side does what, and exactly one of them is true: + // `replicateRightSide` is the negation of `replicateLeftSide`, and the branch above is taken + // only when the side it picked may replicate, so `applyPartialClustering` holds with one flag + // set. That split is what makes an aligned pair of these layouts sound, and it is the whole + // of it: for a key, a partition holding part of it on the spread side holds all of it on the + // side that repeats the group, so pairing the two index by index loses no match. + // `ValidateRequirements` reads such a pair as aligned and cannot tell it from two sides that + // split the key between them, so the decision has to be right here. ( GroupPartitionsExec(rawLeft, leftSpec.joinKeyPositions, Some(mergedPartitionKeys), leftReducers, diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/ValidateRequirements.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/ValidateRequirements.scala index 1ac6b809fd25..6816e3a05622 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/ValidateRequirements.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/ValidateRequirements.scala @@ -21,6 +21,8 @@ import org.apache.spark.internal.Logging import org.apache.spark.sql.catalyst.expressions._ import org.apache.spark.sql.catalyst.plans.physical._ import org.apache.spark.sql.execution._ +import org.apache.spark.sql.execution.joins.ShuffledJoin +import org.apache.spark.sql.internal.SQLConf /** * Validates that the [[org.apache.spark.sql.catalyst.plans.physical.Partitioning Partitioning]] @@ -45,29 +47,81 @@ object ValidateRequirements extends Logging { assert(requiredChildDistributions.length == children.length) assert(requiredChildOrderings.length == children.length) + // A `ClusteredDistribution` is the one distribution an operator can owe its children together + // rather than one by one, so an operator whose children all owe one is judged on their mutual + // layout below, and every other child, an operator with a single clustered child included, + // answers for itself. That is every such operator, not only a join: one that zips corresponding + // partitions, a cogroup for instance, reads a layout both children have to hold together. + val clusteredMultiChild = children.length > 1 && + requiredChildDistributions.forall(_.isInstanceOf[ClusteredDistribution]) + + // The one member a finished plan may report without satisfying the distribution is the shape + // partially clustered distribution spreads ungrouped, and one producer builds it: + // `EnsureRequirements.checkKeyGroupCompatible`. Both halves of that admission are asked here, + // and the answer is passed down to the pairing below so the waiver cannot be read one way here + // and the other way there. Neither half is a second copy: the operator kinds come from the + // producer itself. What is left to the member, its count and the permission for the collapse it + // went through, is asked there. + val mayBeUngrouped = clusteredMultiChild && + SQLConf.get.v2BucketingPartiallyClusteredDistributionEnabled && + ShuffledJoin.partiallyClusteredJoinType(plan).isDefined + val satisfied = children.zip(requiredChildDistributions.zip(requiredChildOrderings)).forall { case (child, (distribution, ordering)) - if !child.outputPartitioning.satisfies(distribution) + if (!child.outputPartitioning.satisfies(distribution) && + !(mayBeUngrouped && + PartitioningCollection.representativeOf(child.outputPartitioning).isDefined)) || !SortOrder.orderingSatisfies(child.outputOrdering, ordering) => logDebug(s"ValidateRequirements failed: $distribution, $ordering\n$plan") false case _ => true } - if (satisfied && children.length > 1 && - requiredChildDistributions.forall(_.isInstanceOf[ClusteredDistribution])) { - // Check the co-partitioning requirement. - val specs = children.map(_.outputPartitioning).zip(requiredChildDistributions).map { - case (p, d) => p.createShuffleSpec(d.asInstanceOf[ClusteredDistribution]) - } - if (specs.tail.forall(_.isCompatibleWith(specs.head))) { - true - } else { + // What a multi-child clustered operator reads is the pairing: a pair aligned without grouping, + // which partially clustered distribution builds on purpose, is one the sides agree on while + // neither is grouped. The pairing cannot tell how the two sides hold a key's rows, since a + // spread side and one that repeats the whole group report the same keys as two sides that + // split the key, so that rests on the producer, which is why the ungrouped shape alone is + // waived above. + if (!satisfied) { + false + } else if (clusteredMultiChild) { + val paired = satisfiesForPairing(children, requiredChildDistributions, mayBeUngrouped) + if (!paired) { logDebug(s"ValidateRequirements failed: children not co-partitioned in\n$plan") - false } + paired } else { - satisfied + true + } + } + + /** + * Whether the sides of a multi-child clustered operator line up: every side offers the layouts it + * reports ([[PartitioningCollection.specsForPairing]]), and one member of the first side pairs + * with every other side. A plan holds what its members report, so no key is deduped and none is + * re-sorted to make a pair: a side is judged on the partitions it has, under the key the + * operation clusters on. + * + * This is the question `EnsureRequirements` asks of a pair it takes as it stands, the + * `compatibleAsIs` path. Its other path commits on the sides it builds rather than on the pair it + * picked: `agreeingPairs` admits a pair a reduce would reconcile (`areKeysCompatible` with the + * reduce allowed), and `committed` then compares the declared layouts of the two sides that + * reduce ran on. A finished plan is asked the strict question alone, which is sound because the + * reduced pair answers it: `isExpressionCompatible` reads its two sides through + * `hasSameReducedKeys`. The coverage of every operation key a member is additionally asked for + * (`spark.sql.requireAllClusterKeysForCoPartition`) is a skew heuristic, and is not part of it. + */ + private def satisfiesForPairing( + children: Seq[SparkPlan], + distributions: Seq[Distribution], + mayBeUngrouped: Boolean): Boolean = { + val specs = children.zip(distributions).map { case (child, distribution) => + PartitioningCollection.specsForPairing( + child.outputPartitioning, distribution.asInstanceOf[ClusteredDistribution], mayBeUngrouped) + } + specs.headOption.exists { firstSide => + firstSide.exists(head => specs.tail.forall(side => side.exists(_.isCompatibleWith(head)))) } } } diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/joins/ShuffledJoin.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/joins/ShuffledJoin.scala index 7d9d7adc15e7..9ebf8c9aee4a 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/joins/ShuffledJoin.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/joins/ShuffledJoin.scala @@ -20,6 +20,7 @@ package org.apache.spark.sql.execution.joins import org.apache.spark.sql.catalyst.expressions.{Attribute, Expression} import org.apache.spark.sql.catalyst.plans.{ExistenceJoin, FullOuter, InnerLike, JoinType, LeftAnti, LeftExistence, LeftOuter, LeftSingle, RightOuter} import org.apache.spark.sql.catalyst.plans.physical.{ClusteredDistribution, Distribution, KeyedPartitioning, Partitioning, PartitioningCollection, UnknownPartitioning, UnspecifiedDistribution} +import org.apache.spark.sql.execution.SparkPlan import org.apache.spark.sql.internal.SQLConf /** @@ -173,4 +174,20 @@ object ShuffledJoin { case _: InnerLike | RightOuter => true case _ => false } + + /** + * The join type of the operators whose partially clustered alignment spreads one side against a + * repeater, and `None` for every other operator. `EnsureRequirements.checkKeyGroupCompatible` is + * the only producer of such a pair, and `ValidateRequirements` waives an ungrouped side only for + * the operators this names, so the two read one list rather than a copy each. + * + * A `SortMergeAsOfJoinExec` is a `ShuffledJoin` and builds none, which is why the kinds are named + * rather than read off the trait. The side that repeats is decided next to this, by + * `canDuplicateLeftSide` and `canDuplicateRightSide`. + */ + def partiallyClusteredJoinType(plan: SparkPlan): Option[JoinType] = plan match { + case smj: SortMergeJoinExec => Some(smj.joinType) + case sj: ShuffledHashJoinExec => Some(sj.joinType) + case _ => None + } } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala index 6cb17713e1c8..44065c3d0ccb 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala @@ -27,7 +27,7 @@ import org.apache.spark.sql.catalyst.expressions.{Ascending, AttributeReference, import org.apache.spark.sql.catalyst.expressions.aggregate.Complete import org.apache.spark.sql.catalyst.plans.{Cross, ExistenceJoin, Inner, JoinType, LeftAnti, LeftSemi, LeftSingle} import org.apache.spark.sql.catalyst.plans.physical -import org.apache.spark.sql.catalyst.plans.physical.KeyedPartitioning +import org.apache.spark.sql.catalyst.plans.physical.{KeyedPartitioning, PartitioningCollection} import org.apache.spark.sql.connector.catalog.{Column, Identifier, InMemoryCatalystRuntimeFilterCatalog, InMemoryTableCatalog} import org.apache.spark.sql.connector.catalog.functions._ import org.apache.spark.sql.connector.distributions.Distributions @@ -45,6 +45,7 @@ import org.apache.spark.sql.execution.{ SparkPlan, UnionExec, WholeStageCodegenExec} +import org.apache.spark.sql.execution.adaptive.{AQEShuffleReadExec, ResultQueryStageExec} import org.apache.spark.sql.execution.aggregate.{BaseAggregateExec, HashAggregateExec, SortAggregateExec} import org.apache.spark.sql.execution.datasources.v2.{BatchScanExec, DataSourceV2ScanRelation, GroupPartitionsExec} import org.apache.spark.sql.execution.exchange.{EnsureRequirements, ReusedExchangeExec, ShuffleExchangeExec, ShuffleExchangeLike, ValidateRequirements} @@ -9228,6 +9229,101 @@ class KeyGroupedPartitioningSuite } } + test("SPARK-59671: a partially clustered join leaves AQE's shuffle coalescing alone") { + // AQE validates a stage's whole candidate plan before accepting a shuffle-read change: on the + // base, a storage-partitioned join whose sides are aligned but not grouped kept every shuffle + // in its stage uncoalesced, unrelated ones included. Partially clustered distribution plans + // such a pair: the side that keeps its splits spreads them, and the other replicates its + // group across them, so both report repeated keys on purpose. The assertions below pin the + // read that the join's presence must leave alone. + val idCols = Array(Column.create("id", IntegerType), Column.create("data", StringType)) + createTable("pc1", idCols, Array(identity("id"))) + createTable("pc2", idCols, Array(identity("id"))) + sql("INSERT INTO testcat.ns.pc1 VALUES (1, 'a1'), (2, 'a2'), (3, 'a3')") + // Key 1 twice: the side holding it keeps and spreads its two splits, which is what makes the + // pair ungrouped while its keys still line up. + sql("INSERT INTO testcat.ns.pc2 VALUES (1, 'b1'), (1, 'b1b'), (2, 'b2'), (4, 'b4')") + + withSQLConf( + SQLConf.V2_BUCKETING_ENABLED.key -> "true", + SQLConf.V2_BUCKETING_PUSH_PART_VALUES_ENABLED.key -> "true", + SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> "true", + "spark.sql.autoBroadcastJoinThreshold" -> "-1") { + val df = sql( + s""" + |SELECT /*+ MERGE(a) */ a.id, b.data + |FROM testcat.ns.pc1 a JOIN testcat.ns.pc2 b ON a.id = b.id + |UNION ALL + |SELECT count(*), cast(id % 2 AS STRING) + |FROM testcat.ns.pc1 GROUP BY id % 2 + |""".stripMargin) + checkAnswer(df, Seq(Row(1, "b1"), Row(1, "b1b"), Row(2, "b2"), Row(2, "1"), Row(1, "0"))) + + val plan = stripAQEPlan(df.queryExecution.executedPlan) + val joins = collect(plan) { case j: ShuffledJoin => j } + assert(joins.size == 1, s"test setup: one storage-partitioned join:\n$plan") + assert(collectGroupPartitions(plan).exists { g => + PartitioningCollection.representativeOf(g.outputPartitioning).exists(!_.isGrouped) + }, s"test setup: the pair is spread, so a side repeats its keys:\n$plan") + // The join's side shuffles nothing (this suite's `collectShuffles` counts the exchanges a + // join reads through), and the chain shuffles once, for the aggregate: the stage holding + // both is the one whose read the join's presence could have kept uncoalesced. + assert(collectShuffles(plan).isEmpty, s"the join shuffles nothing:\n$plan") + assert(collect(plan) { case s: ShuffleExchangeExec => s }.size === 1, + s"and the chain shuffles once, for the aggregate:\n$plan") + val aqeReads = collect(df.queryExecution.executedPlan) { case r: AQEShuffleReadExec => r } + assert(aqeReads.size === 1 && aqeReads.head.hasCoalescedPartition, + s"the aggregate's shuffle read must coalesce:\n${df.queryExecution.executedPlan}") + + // And the two share a stage, which is what makes the coalesce a decision the join can + // block: a read in a stage of its own would coalesce whatever the join did. + val finalStage = collect(df.queryExecution.executedPlan) { + case s: ResultQueryStageExec => s + } + assert(finalStage.size === 1 && finalStage.head.plan.exists { + case j: ShuffledJoin => true + case _ => false + }, s"test setup: the join is in the stage the read belongs to:\n" + + s"${df.queryExecution.executedPlan}") + } + } + + test("SPARK-59671: a three-table chain keeps its AQE coalescing") { + val idCols = Array(Column.create("id", IntegerType), Column.create("data", StringType)) + createTable("p3a", idCols, Array(identity("id"))) + createTable("p3b", idCols, Array(identity("id"))) + createTable("p3c", idCols, Array(identity("id"))) + sql("INSERT INTO testcat.ns.p3a VALUES (1, 'a1'), (2, 'a2')") + sql("INSERT INTO testcat.ns.p3b VALUES (1, 'b1'), (1, 'b1b'), (2, 'b2')") + sql("INSERT INTO testcat.ns.p3c VALUES (1, 'c1'), (2, 'c2')") + withSQLConf( + SQLConf.V2_BUCKETING_ENABLED.key -> "true", + SQLConf.V2_BUCKETING_PUSH_PART_VALUES_ENABLED.key -> "true", + SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> "true", + "spark.sql.autoBroadcastJoinThreshold" -> "-1") { + val df = sql( + s""" + |SELECT /*+ MERGE(a, b), MERGE(a, c) */ a.id AS aid, b.id AS bid, c.data + |FROM testcat.ns.p3a a JOIN testcat.ns.p3b b ON a.id = b.id + |JOIN testcat.ns.p3c c ON a.id = c.id + |UNION ALL + |SELECT count(*), 0, 'x' FROM testcat.ns.p3a GROUP BY id % 2 + |""".stripMargin) + val expected = Seq(Row(1, 1, "c1"), Row(1, 1, "c1"), Row(2, 2, "c2"), + Row(1, 0, "x"), Row(1, 0, "x")) + assert(df.collect().map(_.toString).sorted === expected.map(_.toString).sorted, + s"rows must come out whole:\n${df.queryExecution.executedPlan}") + val plan = stripAQEPlan(df.queryExecution.executedPlan) + assert(collectShuffles(plan).isEmpty, s"the chain must not shuffle:\n$plan") + assert(plan.exists(p => PartitioningCollection.flatten(p.outputPartitioning) + .count(_.isInstanceOf[KeyedPartitioning]) >= 2), + s"test setup: a side reports one keyed member per join key column:\n$plan") + assert(ValidateRequirements.validate(plan), s"the chain's pairing holds up:\n$plan") + assert(collect(df.queryExecution.executedPlan) { case r: AQEShuffleReadExec => r } + .exists(_.hasCoalescedPartition), + s"the aggregate's shuffle read must coalesce:\n${df.queryExecution.executedPlan}") + } + } } /** diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/exchange/ValidateRequirementsSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/exchange/ValidateRequirementsSuite.scala index 4e02a10eb41e..9ef546aefe07 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/exchange/ValidateRequirementsSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/exchange/ValidateRequirementsSuite.scala @@ -17,12 +17,17 @@ package org.apache.spark.sql.execution.exchange -import org.apache.spark.sql.catalyst.expressions.{Ascending, SortOrder} +import org.apache.spark.sql.catalyst.InternalRow +import org.apache.spark.sql.catalyst.expressions.{Ascending, AttributeReference, GreaterThan, Literal, SortOrder} +import org.apache.spark.sql.catalyst.optimizer.BuildLeft import org.apache.spark.sql.catalyst.plans.Inner -import org.apache.spark.sql.catalyst.plans.physical.{HashPartitioning, SinglePartition} -import org.apache.spark.sql.execution.SortExec -import org.apache.spark.sql.execution.joins.SortMergeJoinExec +import org.apache.spark.sql.catalyst.plans.physical.{ClusteredDistribution, HashPartitioning, KeyedPartitioning, PartitioningCollection, SinglePartition} +import org.apache.spark.sql.execution.{CoGroupExec, DummySparkPlan, SortExec, SparkPlan} +import org.apache.spark.sql.execution.datasources.v2.GroupPartitionsExec +import org.apache.spark.sql.execution.joins.{ShuffledHashJoinExec, SortMergeAsOfJoinExec, SortMergeJoinExec} +import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.test.SharedSparkSession +import org.apache.spark.sql.types.{IntegerType, ObjectType} class ValidateRequirementsSuite extends SharedSparkSession { @@ -158,4 +163,336 @@ class ValidateRequirementsSuite extends SharedSparkSession { testNestedJoin(Seq((2, 2), (1, 1)), Seq((2, 2)), Seq(5, 5, 5), false) testNestedJoin(Seq((2, 2), (1, 1)), Seq((2, 5)), Seq(5, 5, 5), false) } + + test("SPARK-59671: a co-partitioning operator judges keyed children by their pairing") { + // The sides of a storage-partitioned join aligned for skew repeat their spread keys on + // purpose: neither satisfies a clustering on its own, yet the two key sequences agree + // index by index, and the pairing is what the operator reads. + val a = AttributeReference("a", IntegerType)() + val b = AttributeReference("b", IntegerType)() + val rows = Seq(InternalRow(1), InternalRow(1), InternalRow(2)) + val left = DummySparkPlan(outputPartitioning = KeyedPartitioning(Seq(a), rows)) + val right = DummySparkPlan(outputPartitioning = KeyedPartitioning(Seq(b), rows)) + val join = ShuffledHashJoinExec(Seq(a), Seq(b), Inner, BuildLeft, None, left, right) + withSQLConf(SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> "true") { + assert(ValidateRequirements.validate(join), + s"aligned but ungrouped keyed sides pair, and that pairing is the requirement:\n$join") + + // The same key rows in another order do not pair: position by position is the contract. + val off = DummySparkPlan(outputPartitioning = + KeyedPartitioning(Seq(b), Seq(InternalRow(1), InternalRow(2), InternalRow(1)))) + assert(!ValidateRequirements.validate(join.copy(right = off)), + "the same keys in a different order are not aligned") + + // Nor does a keyed side pair with one that never pairs: a hashed side matches neither the + // keys nor the layout of a keyed one. + val hashed = DummySparkPlan(outputPartitioning = HashPartitioning(Seq(b), 3)) + assert(!ValidateRequirements.validate(join.copy(right = hashed)), + "a keyed side does not pair with a hashed one") + } + + // An ungrouped side is a plan only partially clustered distribution builds, so with that off + // the per-side check stands: nothing is left to group the side, and a pair the plan does not + // hold is not one the operator reads. + assert(!ValidateRequirements.validate(join), + s"an ungrouped side is admitted only where something builds one:\n$join") + + // And a pair each side satisfies on its own is still refused when the sides do not line up, + // which is the operator's requirement: the pairing, not the per-side answer. + val leftKeys = DummySparkPlan(outputPartitioning = + KeyedPartitioning(Seq(a), Seq(InternalRow(1), InternalRow(2), InternalRow(3)))) + val rightKeys = DummySparkPlan(outputPartitioning = + KeyedPartitioning(Seq(b), Seq(InternalRow(4), InternalRow(5), InternalRow(6)))) + assert(leftKeys.outputPartitioning.satisfies(ClusteredDistribution(Seq(a))) && + rightKeys.outputPartitioning.satisfies(ClusteredDistribution(Seq(b))), + "test setup: each side answers its own clustering") + assert(!ValidateRequirements.validate(ShuffledHashJoinExec( + Seq(a), Seq(b), Inner, BuildLeft, None, leftKeys, rightKeys)), + "the sides do not line up, and nothing else the operator reads says otherwise") + } + + test("SPARK-59671: a keyed pair is judged on the layouts the plan holds") { + // Both permissions on, so the shape is the one partially clustered distribution spreads. + // Neither side's keys are deduped or re-sorted to make the pair line up: reading it through + // `createShuffleSpec` would do both, so two sides whose layouts disagree as they stand would + // answer as though they agreed. + val a = AttributeReference("a", IntegerType)() + val b = AttributeReference("b", IntegerType)() + val rows = Seq(InternalRow(1), InternalRow(1), InternalRow(2)) + withSQLConf( + SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "true", + SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> "true") { + val left = DummySparkPlan(outputPartitioning = KeyedPartitioning(Seq(a), rows)) + // The same keys in another order, which a projection onto distinct sorted keys would + // normalize away. + val off = DummySparkPlan(outputPartitioning = + KeyedPartitioning(Seq(b), Seq(InternalRow(1), InternalRow(2), InternalRow(1)))) + assert(!ValidateRequirements.validate(ShuffledHashJoinExec( + Seq(a), Seq(b), Inner, BuildLeft, None, left, off)), + "the sides report different layouts, and nothing normalizes them") + + // And a count: the other side holds two partitions, and neither side has the node that would + // make the two agree. + val twoKeys = DummySparkPlan(outputPartitioning = + KeyedPartitioning(Seq(b), Seq(InternalRow(1), InternalRow(2)))) + assert(!ValidateRequirements.validate(ShuffledHashJoinExec( + Seq(a), Seq(b), Inner, BuildLeft, None, left, twoKeys)), + "a child holding three partitions does not pair with one holding two") + } + } + + test("SPARK-59671: a pair is judged on the key the operation clusters on") { + // The subset permission applies where the operation's keys are a subset of the source's + // partitioning keys: a side partitioned on `[a, b]` can serve an operator on `[a]`. Such a side + // is judged on its own partitions under `[a]`, which is what its own key already gives it when + // dropping `b` merges no partition, whether or not a node stands over it. No key is deduped and + // none is re-sorted, so the pairs that line up are the ones whose keys line up as reported. + val a = AttributeReference("a", IntegerType)() + val b = AttributeReference("b", IntegerType)() + val c = AttributeReference("c", IntegerType)() + val d = AttributeReference("d", IntegerType)() + withSQLConf( + SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "true", + SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> "false") { + def joinOnKeys(left: SparkPlan, right: SparkPlan): SparkPlan = + ShuffledHashJoinExec(Seq(a), Seq(c), Inner, BuildLeft, None, left, right) + + // The projection's result, which is what such a plan reports from the side it grouped. + val projected = DummySparkPlan(outputPartitioning = + KeyedPartitioning(Seq(a), Seq(InternalRow(1), InternalRow(2)))) + val projectedRight = DummySparkPlan(outputPartitioning = + KeyedPartitioning(Seq(c), Seq(InternalRow(1), InternalRow(2)))) + assert(ValidateRequirements.validate(joinOnKeys(projected, projectedRight)), + "the layout the grouping node leaves is what the pair is judged on") + + // The source's keys, one step earlier: the second expression is the operation's to drop, and + // what is left is the same pair, so it stands. + val source = DummySparkPlan(outputPartitioning = + KeyedPartitioning(Seq(a, b), Seq(InternalRow(1, 1), InternalRow(2, 2)))) + val sourceRight = DummySparkPlan(outputPartitioning = + KeyedPartitioning(Seq(c, d), Seq(InternalRow(1, 1), InternalRow(2, 2)))) + assert(ValidateRequirements.validate(joinOnKeys(source, sourceRight)), + s"the key the operation clusters on is what the side offers:\n$source") + + // And it is offered in the order it is reported, so a side whose keys run the other way is + // not aligned: a spec that sorted them would call these two the same pair. + val reversed = DummySparkPlan(outputPartitioning = + KeyedPartitioning(Seq(a, b), Seq(InternalRow(2, 2), InternalRow(1, 1)))) + assert(!ValidateRequirements.validate(joinOnKeys(reversed, sourceRight)), + s"the keys are compared as reported, not as sorted:\n$reversed") + } + } + + test("SPARK-59671: a single clustered child still owes its own grouping") { + // The pairing stands in for the per-child check only where children pair with each other. + // An operator with a single clustered child, an aggregate over a join output say, is judged + // per side: the ungrouped keyed layout does not satisfy it until a grouping stands under. + val a = AttributeReference("a", IntegerType)() + val child = DummySparkPlan(outputPartitioning = + KeyedPartitioning(Seq(a), Seq(InternalRow(1), InternalRow(1)))) + val parent = DummySparkPlan( + children = Seq(child), + requiredChildDistribution = Seq(ClusteredDistribution(Seq(a))), + requiredChildOrdering = Seq(Nil)) + assert(!ValidateRequirements.validate(parent), "an ungrouped child fails a lone clustered slot") + } + + test("SPARK-59671: a partially clustered pair planned by the rule passes validation") { + // Partially clustered distribution spreads a side, so its replicate side reports a + // non-grouped layout whose keys repeat on purpose, and the validator's per-side check + // refused such a pair. AQE validates a stage's whole candidate plan before accepting a + // shuffle-read change, so the plan's stage took no coalescing either. The pairing takes + // it: the spread is deliberate, and the two sides agree index by index. + withSQLConf( + SQLConf.V2_BUCKETING_PUSH_PART_VALUES_ENABLED.key -> "true", + SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> "true") { + val a = AttributeReference("a", IntegerType)() + val b = AttributeReference("b", IntegerType)() + val left = DummySparkPlan(outputPartitioning = + KeyedPartitioning(Seq(a), Seq(InternalRow(1), InternalRow(1), InternalRow(2)))) + val right = DummySparkPlan(outputPartitioning = + KeyedPartitioning(Seq(b), Seq(InternalRow(1), InternalRow(2)))) + val plan = new EnsureRequirements().apply( + SortMergeJoinExec(Seq(a), Seq(b), Inner, None, left, right)) + val join = plan.collectFirst { case j: SortMergeJoinExec => j } + .getOrElse(fail(s"expected the join back:\n${plan.treeString}")) + def groupingAt(plan: SparkPlan): Option[GroupPartitionsExec] = plan match { + case g: GroupPartitionsExec => Some(g) + case s: SortExec if !s.global => groupingAt(s.child) + case _ => None + } + val sideGroupings = join.children.flatMap(groupingAt) + assert(sideGroupings.size == 2, + s"test setup: both sides are aligned:\n${plan.treeString}") + assert(sideGroupings.exists { g => + PartitioningCollection.representativeOf(g.outputPartitioning).exists(!_.isGrouped) + }, s"test setup: a side keeps its splits, so its keys repeat:\n${plan.treeString}") + assert(ValidateRequirements.validate(plan), + s"a spread pair that agrees on its keys is accepted:\n${plan.treeString}") + } + } + + test("SPARK-59671: a collection of keyed members is judged by its pairing") { + // A side that reports several keyed alternatives (a projection over a join keeps one per join + // key column) is judged on the members the admission keeps, not on the collection's own spec + // build having to leave one behind: a side whose members all repeat their keys is a pair the + // planner aligns without grouping, and it holds up when the member keyed on the join key is + // the one that lines up with the other side. + val a = AttributeReference("a", IntegerType)() + val b = AttributeReference("b", IntegerType)() + val rows = Seq(InternalRow(1), InternalRow(1), InternalRow(2)) + def keyed(attr: AttributeReference): DummySparkPlan = + DummySparkPlan(outputPartitioning = KeyedPartitioning(Seq(attr), rows)) + val bothAlternatives = DummySparkPlan(outputPartitioning = + PartitioningCollection.fromPartitionings(Seq( + KeyedPartitioning(Seq(a), rows), KeyedPartitioning(Seq(b), rows)))) + // The alternatives repeat their keys, so they are the layout partially clustered distribution + // spreads, and that is where the admission keeps an ungrouped member. + withSQLConf(SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> "true") { + assert(ValidateRequirements.validate(ShuffledHashJoinExec( + Seq(a), Seq(a), Inner, BuildLeft, None, bothAlternatives, keyed(a))), + "the member keyed on the join key pairs with the other side") + + // The pairing still has to be there: a side whose members are keyed on something else offers + // nothing to pair with. + val wrongKeys = DummySparkPlan(outputPartitioning = + PartitioningCollection.fromPartitionings(Seq( + KeyedPartitioning(Seq(b), rows), KeyedPartitioning(Seq(b), rows)))) + assert(!ValidateRequirements.validate(ShuffledHashJoinExec( + Seq(a), Seq(a), Inner, BuildLeft, None, wrongKeys, keyed(a))), + "a side offering no member keyed on the join keys does not pair") + } + } + + test("SPARK-59671: a multi-child clustered operator is judged on its children's pairing") { + // Every operator whose children all owe a `ClusteredDistribution` reads one layout they have to + // hold together, whichever operator it is: a cogroup zips corresponding partitions, so two + // sides that each satisfy the distribution on their own are not enough. + val a = AttributeReference("a", IntegerType)() + val b = AttributeReference("b", IntegerType)() + val rows = Seq(InternalRow(1), InternalRow(1), InternalRow(2)) + def cogroupOf(left: SparkPlan, right: SparkPlan): SparkPlan = CoGroupExec( + (key: Any, l: Iterator[Any], r: Iterator[Any]) => Nil, + Literal(1), Literal(1), Literal(1), Seq(a), Seq(b), Seq(a), Seq(b), Nil, Nil, + AttributeReference("obj", ObjectType(classOf[AnyRef]))(), left, right) + def grouped(attr: AttributeReference, keys: Seq[Int]): DummySparkPlan = DummySparkPlan( + outputOrdering = Seq(SortOrder(attr, Ascending)), + outputPartitioning = KeyedPartitioning(Seq(attr), keys.map(InternalRow(_)))) + + withSQLConf(SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> "true") { + // The ungrouped shape is waived only for the producer that spreads a side, so a cogroup over + // the very pair a join reads is refused: it would run its function once per spread part. + val left = DummySparkPlan( + outputOrdering = Seq(SortOrder(a, Ascending)), + outputPartitioning = KeyedPartitioning(Seq(a), rows)) + val right = DummySparkPlan( + outputOrdering = Seq(SortOrder(b, Ascending)), + outputPartitioning = KeyedPartitioning(Seq(b), rows)) + assert(ValidateRequirements.validate( + ShuffledHashJoinExec(Seq(a), Seq(b), Inner, BuildLeft, None, left, right)), + "test setup: the join reads this pair") + assert(!ValidateRequirements.validate(cogroupOf(left, right)), + "a cogroup reads no spread side, and this pair passes nothing else") + } + + // What a cogroup shares with a join is the mutual check on the layouts its sides report. + def cogroup(aKeys: Seq[Int], bKeys: Seq[Int]): SparkPlan = + cogroupOf(grouped(a, aKeys), grouped(b, bKeys)) + assert(ValidateRequirements.validate(cogroup(Seq(1, 2), Seq(1, 2))), + "two sides holding the same grouped layout are read as they stand") + assert(!ValidateRequirements.validate(cogroup(Seq(1, 2), Seq(1, 2, 3))), + "a side holding three partitions does not pair with one holding two") + assert(!ValidateRequirements.validate(cogroup(Seq(1, 2), Seq(2, 1))), + "the same keys in another order are not aligned") + } + + test("SPARK-59671: an ungrouped pair is waived only where the producer spreads a side") { + // `EnsureRequirements.checkKeyGroupCompatible` is the path that spreads one side against a + // repeater, and it is entered for a sort-merge or shuffled-hash join. A sort-merge as-of join + // is a `ShuffledJoin` too and builds none, so a pair it reads ungrouped is one no producer made + // and the per-side check stands. + val a = AttributeReference("a", IntegerType)() + val b = AttributeReference("b", IntegerType)() + val rows = Seq(InternalRow(1), InternalRow(1), InternalRow(2)) + withSQLConf(SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> "true") { + val left = DummySparkPlan( + outputOrdering = Seq(SortOrder(a, Ascending), SortOrder(a, Ascending)), + outputPartitioning = KeyedPartitioning(Seq(a), rows)) + val right = DummySparkPlan( + outputOrdering = Seq(SortOrder(b, Ascending), SortOrder(b, Ascending)), + outputPartitioning = KeyedPartitioning(Seq(b), rows)) + assert(ValidateRequirements.validate( + ShuffledHashJoinExec(Seq(a), Seq(b), Inner, BuildLeft, None, left, right)), + "test setup: the join the waiver is for reads this pair") + assert(!ValidateRequirements.validate(SortMergeAsOfJoinExec( + Seq(a), Seq(b), Seq(a), Seq(b), GreaterThan(a, b), a, Inner, None, left, right)), + "an as-of join builds no spread side, so the same pair is refused") + } + } + + test("SPARK-59671: a side is judged on the member that pairs, whichever one it is") { + // A side that reports several keyed alternatives offers all of them: the one keyed on the + // join key is the one that lines up, and it need not be the first the side reports. Reading a + // single member would refuse this pair, which the planner builds whenever a projection keeps + // two key columns. + val a = AttributeReference("a", IntegerType)() + val b = AttributeReference("b", IntegerType)() + val rows = Seq(InternalRow(1), InternalRow(1), InternalRow(2)) + def keyed(attr: AttributeReference): KeyedPartitioning = KeyedPartitioning(Seq(attr), rows) + val left = DummySparkPlan(outputPartitioning = + PartitioningCollection.fromPartitionings(Seq(keyed(b), keyed(a)))) + val right = DummySparkPlan(outputPartitioning = keyed(a)) + withSQLConf(SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> "true") { + assert(ValidateRequirements.validate(ShuffledHashJoinExec( + Seq(a), Seq(a), Inner, BuildLeft, None, left, right)), + "the second member is the one keyed on the join key, and it pairs") + } + } + + test("SPARK-59671: the ordering requirement is not part of the exemption") { + // A pair can line up and still owe its operator an ordering: the exemption covers the + // distribution clause alone, so a sort-merge join over the aligned sides, with nothing + // ordering them, is refused. + val a = AttributeReference("a", IntegerType)() + val b = AttributeReference("b", IntegerType)() + val rows = Seq(InternalRow(1), InternalRow(1), InternalRow(2)) + val left = DummySparkPlan(outputPartitioning = KeyedPartitioning(Seq(a), rows)) + val right = DummySparkPlan(outputPartitioning = KeyedPartitioning(Seq(b), rows)) + assert(left.outputOrdering.isEmpty && right.outputOrdering.isEmpty, + "test setup: nothing orders the sides") + withSQLConf(SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> "true") { + assert(!ValidateRequirements.validate( + SortMergeJoinExec(Seq(a), Seq(b), Inner, None, left, right)), + "the sides pair, and still owe the join its ordering") + } + } + + test("SPARK-59671: a collapsed pair is not admitted on its pairing alone") { + // A layout whose keys were collapsed (a partition standing for several of the source's) serves + // a clustering through a grouping node only where that grouping is permitted, and the + // permission is a config. Without it nothing admits such a member, and a side offering nothing + // has nothing to pair with. + val a = AttributeReference("a", IntegerType)() + val b = AttributeReference("b", IntegerType)() + def collapsed(attr: AttributeReference): DummySparkPlan = DummySparkPlan( + outputPartitioning = KeyedPartitioning(Seq(attr), + Seq(InternalRow(1), InternalRow(1), InternalRow(2))) + .withLayout(_.copy(isCollapsed = true))) + def pair: SparkPlan = ShuffledHashJoinExec(Seq(a), Seq(b), Inner, BuildLeft, None, + collapsed(a), collapsed(b)) + + withSQLConf(SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "false") { + assert(!ValidateRequirements.validate(pair), + s"a collapsed pair is refused while the subset permission is off:\n$pair") + } + // And where an ungrouped member is admitted at all, which takes partially clustered + // distribution, the collapse still has to be permitted: the waiver carries the producer's + // shape, not a permission the producer would not have had. + withSQLConf( + SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "false", + SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> "true") { + assert(!ValidateRequirements.validate(pair), + s"the ungrouped waiver does not carry a collapse the planner would not group:\n$pair") + } + } } From 32ddb2663e4d89f15b26d68fcbd1f306b3e700f8 Mon Sep 17 00:00:00 2001 From: Xiduo You Date: Thu, 24 Sep 2026 09:31:33 +0800 Subject: [PATCH 2/2] [SPARK-59671][SQL] Narrow the ungrouped waiver to join types a side may be duplicated for The waiver let through any sort-merge or shuffled-hash join, but the producer spreads a side only where the join type may duplicate one: `EnsureRequirements` clears the replicate side it picked when `canDuplicateLeftSide` / `canDuplicateRightSide` rejects the join type and sets no partially clustered distribution, so no side of a `FullOuter` join is ever spread. `ValidateRequirements` asks that second gate now, at the waiver rather than through `ShuffledJoin.partiallyClusteredJoinType`, which is the producer's entry to key-group checking altogether: a kind turned away there would lose the storage-partitioned join it can still plan without spreading a side. `ShuffledJoin.partiallyClusteredJoinType` says as much in its doc, and `satisfiesForPairing` describes its relationship to the planner's `compatibleAsIs` path without claiming the same predicate: the planner asks that of two unprojected specs, while the validator may read a member relabelled onto the key the operation clusters on (`reportedSpecOf`). The producer-boundary test pins the join type as well as the operator kind: the ungrouped pair it already refused for an as-of join is refused for a `FullOuter` join, and the `Inner` pair the waiver is for still passes. `ShuffleSpecSuite` and `DistributionSuite` (52 tests) and `ValidateRequirementsSuite`, `EnsureRequirementsSuite`, `GroupPartitionsExecSuite`, `KeyGroupedPartitioningSuite` and `ProjectedOrderingAndPartitioningSuite` (362 tests) are green, with `dev/lint-scala` clean. Assisted-by: DeepSeek Flash --- .../exchange/ValidateRequirements.scala | 20 ++++++++++++++----- .../sql/execution/joins/ShuffledJoin.scala | 4 +++- .../exchange/ValidateRequirementsSuite.scala | 7 ++++++- 3 files changed, 24 insertions(+), 7 deletions(-) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/ValidateRequirements.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/ValidateRequirements.scala index 6816e3a05622..aa33a422e5a9 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/ValidateRequirements.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/ValidateRequirements.scala @@ -60,11 +60,19 @@ object ValidateRequirements extends Logging { // `EnsureRequirements.checkKeyGroupCompatible`. Both halves of that admission are asked here, // and the answer is passed down to the pairing below so the waiver cannot be read one way here // and the other way there. Neither half is a second copy: the operator kinds come from the - // producer itself. What is left to the member, its count and the permission for the collapse it - // went through, is asked there. + // producer, and so does the join type's capability, the second gate it applies before spreading + // a side (`canDuplicateLeftSide` / `canDuplicateRightSide`; a `FullOuter` join may duplicate + // neither side, so nothing spreads such a pair). That gate is asked here rather than through + // `partiallyClusteredJoinType`, which is the producer's entry to key-group checking at all: a + // kind turned away there would lose the storage-partitioned join it can still plan without + // spreading. What is left to the member, its count and the permission for the collapse it went + // through, is asked in the spec. val mayBeUngrouped = clusteredMultiChild && SQLConf.get.v2BucketingPartiallyClusteredDistributionEnabled && - ShuffledJoin.partiallyClusteredJoinType(plan).isDefined + ShuffledJoin.partiallyClusteredJoinType(plan).exists { joinType => + ShuffledJoin.canDuplicateLeftSide(joinType) || + ShuffledJoin.canDuplicateRightSide(joinType) + } val satisfied = children.zip(requiredChildDistributions.zip(requiredChildOrderings)).forall { case (child, (distribution, ordering)) @@ -103,8 +111,10 @@ object ValidateRequirements extends Logging { * re-sorted to make a pair: a side is judged on the partitions it has, under the key the * operation clusters on. * - * This is the question `EnsureRequirements` asks of a pair it takes as it stands, the - * `compatibleAsIs` path. Its other path commits on the sides it builds rather than on the pair it + * That is the question `EnsureRequirements` asks of a pair it takes as it stands, though not by + * the same predicate: its `compatibleAsIs` path reads two unprojected specs, while a member here + * may be relabelled onto the key the operation clusters on (see `reportedSpecOf`). Its other path + * commits on the sides it builds rather than on the pair it * picked: `agreeingPairs` admits a pair a reduce would reconcile (`areKeysCompatible` with the * reduce allowed), and `committed` then compares the declared layouts of the two sides that * reduce ran on. A finished plan is asked the strict question alone, which is sound because the diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/joins/ShuffledJoin.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/joins/ShuffledJoin.scala index 9ebf8c9aee4a..9194e5a8bb5a 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/joins/ShuffledJoin.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/joins/ShuffledJoin.scala @@ -182,7 +182,9 @@ object ShuffledJoin { * the operators this names, so the two read one list rather than a copy each. * * A `SortMergeAsOfJoinExec` is a `ShuffledJoin` and builds none, which is why the kinds are named - * rather than read off the trait. The side that repeats is decided next to this, by + * rather than read off the trait. One of these kinds can still be unable to spread for its join + * type, which is a separate question: no side of a `FullOuter` join may be duplicated, so the + * alignment spreads nothing there. The side that repeats is decided next to this, by * `canDuplicateLeftSide` and `canDuplicateRightSide`. */ def partiallyClusteredJoinType(plan: SparkPlan): Option[JoinType] = plan match { diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/exchange/ValidateRequirementsSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/exchange/ValidateRequirementsSuite.scala index 9ef546aefe07..c78a1777de52 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/exchange/ValidateRequirementsSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/exchange/ValidateRequirementsSuite.scala @@ -20,7 +20,7 @@ package org.apache.spark.sql.execution.exchange import org.apache.spark.sql.catalyst.InternalRow import org.apache.spark.sql.catalyst.expressions.{Ascending, AttributeReference, GreaterThan, Literal, SortOrder} import org.apache.spark.sql.catalyst.optimizer.BuildLeft -import org.apache.spark.sql.catalyst.plans.Inner +import org.apache.spark.sql.catalyst.plans.{FullOuter, Inner} import org.apache.spark.sql.catalyst.plans.physical.{ClusteredDistribution, HashPartitioning, KeyedPartitioning, PartitioningCollection, SinglePartition} import org.apache.spark.sql.execution.{CoGroupExec, DummySparkPlan, SortExec, SparkPlan} import org.apache.spark.sql.execution.datasources.v2.GroupPartitionsExec @@ -427,6 +427,11 @@ class ValidateRequirementsSuite extends SharedSparkSession { assert(!ValidateRequirements.validate(SortMergeAsOfJoinExec( Seq(a), Seq(b), Seq(a), Seq(b), GreaterThan(a, b), a, Inner, None, left, right)), "an as-of join builds no spread side, so the same pair is refused") + // Nor for a join type that may duplicate neither side: the producer applies that gate before + // it spreads, so a full outer join over the same pair is one no producer builds either. + assert(!ValidateRequirements.validate( + SortMergeJoinExec(Seq(a), Seq(b), FullOuter, None, left, right)), + "neither side of a full outer join may be duplicated, so nothing spreads this pair") } }