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 3168f552f4c33..1d41624112a90 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 @@ -936,9 +936,13 @@ object KeyedPartitioning { * `PartitioningCollection`, whose invariant requires equal partition keys -- but join types that * expose only one side's partitioning (e.g. LEFT OUTER) run nothing that compares the two * orders, and silently return wrong results. + * + * It is the keys' own ordering, the one `InternalRowComparableWrapper.equals` compares with, so + * one definition answers both. `EnsureRequirements`' `OrderedDistribution` arm is the one place + * that lays grouped keys out in another order, the distribution's own. */ def groupedKeyRowOrdering(dataTypes: Seq[DataType]): BaseOrdering = - RowOrdering.createNaturalAscendingOrdering(dataTypes) + InternalRowComparableWrapper.getInternalRowComparableWrapperFactory(dataTypes).ordering /** * Projects a sequence of partition keys by selecting only the specified positions. diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/util/InternalRowComparableWrapper.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/util/InternalRowComparableWrapper.scala index 5eaac05ade156..d6b841f1dafcb 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/util/InternalRowComparableWrapper.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/util/InternalRowComparableWrapper.scala @@ -152,15 +152,16 @@ object InternalRowComparableWrapper { * Builds wrappers over one row schema, holding the cache lookups that schema needs so a caller * does not repeat them per row. * - * `dataTypes` is what the rows it builds compare at, which is `comparableTypes` of what it was - * given. A caller reporting a type list beside those rows takes it from here rather than erasing - * on its own, so the two cannot answer differently. + * It also answers for the schema it settled on. `dataTypes` is what the rows it builds compare + * at, which is `comparableTypes` of what it was given, and `ordering` sorts them. A caller + * reporting a type list beside those rows, or sorting them, takes it from here rather than + * deriving it again, so the two cannot answer differently. */ final class Factory private[InternalRowComparableWrapper] (val dataTypes: Seq[DataType]) extends (InternalRow => InternalRowComparableWrapper) { + val ordering: BaseOrdering = orderingCache.get(dataTypes) private[this] val structType = structTypeCache.get(dataTypes) - private[this] val ordering = orderingCache.get(dataTypes) override def apply(row: InternalRow): InternalRowComparableWrapper = new InternalRowComparableWrapper(row, dataTypes, structType, ordering) diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/util/InternalRowComparableWrapperSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/util/InternalRowComparableWrapperSuite.scala index 2723eeecfeb3b..26692ca6530be 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/util/InternalRowComparableWrapperSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/util/InternalRowComparableWrapperSuite.scala @@ -19,6 +19,7 @@ package org.apache.spark.sql.catalyst.util import org.apache.spark.SparkFunSuite import org.apache.spark.sql.catalyst.InternalRow +import org.apache.spark.sql.catalyst.plans.physical.KeyedPartitioning import org.apache.spark.sql.types._ class InternalRowComparableWrapperSuite extends SparkFunSuite { @@ -81,6 +82,19 @@ class InternalRowComparableWrapperSuite extends SparkFunSuite { assert(InternalRowComparableWrapper.comparableTypes(alreadyErased) eq alreadyErased) } + test("SPARK-59249: the grouped key layout and wrapper equality hold one ordering instance") { + // Identity is the property to assert, because behaviour is not what changes here: both sides + // were already built by the same function, so they already compared the same way. What one + // instance buys is that neither side can later be given a definition the other does not have. + // `InternalRowComparableWrapper.equals` compares its rows with the instance below, and + // `KeyedPartitioning` sorts and groups partition keys with it. The two type lists are built + // separately, so this pins the shared cache as well. + val wrapper = InternalRowComparableWrapper + .getInternalRowComparableWrapperFactory(Seq(IntegerType, LongType))(InternalRow(1, 2L)) + + assert(KeyedPartitioning.groupedKeyRowOrdering(Seq(IntegerType, LongType)) eq wrapper.ordering) + } + test("SPARK-59187: a factory answers for the types it settled on") { // A caller that reports a type list beside the rows a factory built takes it from here, so the // two cannot answer differently.