diff --git a/sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/SupportsReportOrdering.java b/sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/SupportsReportOrdering.java index 0dc102d11c213..21ed57f890fc9 100644 --- a/sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/SupportsReportOrdering.java +++ b/sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/SupportsReportOrdering.java @@ -35,6 +35,11 @@ public interface SupportsReportOrdering extends Scan { /** * Returns the order in each partition of this data source scan. + *
+ * Spark currently does not rely on a sort order over a transform such as {@code days(ts)} to + * avoid a sort, whether or not it is a partition key. Nor does it rely on the sort orders after + * it. The exception is a sort order on a partition key that is not a transform, when Spark uses + * the reported {@code KeyGroupedPartitioning}. */ SortOrder[] outputOrdering(); } diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala index d8b5ab9c04381..6de00fc795d67 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala @@ -2241,7 +2241,8 @@ object SQLConf { .doc("When enabled, Spark derives output ordering from the partition key expressions of " + "a V2 data source that reports a KeyedPartitioning but does not report explicit ordering " + "via SupportsReportOrdering. Within a single partition all rows share the same key " + - s"value, so the data is trivially sorted by those expressions. Requires " + + "value, so the data is trivially sorted by those expressions. Partition transforms such " + + "as `days(ts)` or `bucket(8, id)` are left out of the ordering. Requires " + s"${V2_BUCKETING_ENABLED.key} to be enabled.") .version("4.2.0") .withBindingPolicy(ConfigBindingPolicy.SESSION) @@ -2253,9 +2254,10 @@ object SQLConf { .doc("When enabled, Spark preserves sort orders over partition key expressions when " + "GroupPartitionsExec coalesces multiple input partitions into one output partition. " + "Because all merged partitions share the same partition key value, sort orders over " + - "those key expressions remain valid after the merge. This applies to both key-derived " + - "ordering (from SupportsReportOrdering) and ordering derived from " + - s"${V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key}. Requires " + + "those key expressions remain valid after the merge. This applies to both the ordering " + + "reported via SupportsReportOrdering and the ordering derived from " + + s"${V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key}. Sort orders over partition " + + "transforms such as `days(ts)` or `bucket(8, id)` are left out. Requires " + s"${V2_BUCKETING_ENABLED.key} to be enabled.") .version("4.2.0") .withBindingPolicy(ConfigBindingPolicy.SESSION) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala index 5c9c38ea72850..d8161c7182c90 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala @@ -19,7 +19,7 @@ package org.apache.spark.sql.execution.datasources.v2 import org.apache.spark.rdd.RDD import org.apache.spark.sql.catalyst.InternalRow -import org.apache.spark.sql.catalyst.expressions.{Ascending, Expression, ExpressionSet, RowOrdering, SortOrder} +import org.apache.spark.sql.catalyst.expressions.{Ascending, Expression, ExpressionSet, RowOrdering, SortOrder, TransformExpression} import org.apache.spark.sql.catalyst.plans.physical import org.apache.spark.sql.catalyst.plans.physical.KeyedPartitioning import org.apache.spark.sql.catalyst.util.truncatedString @@ -116,19 +116,29 @@ trait DataSourceV2ScanExecBase * `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` is on, each partition * contains rows where the key expressions evaluate to a single constant value, so the data * is trivially sorted by those expressions within the partition. + * + * Either way, a sort order that holds a partition transform is dropped, even on a partition key. + * In a reported ordering it also ends the leading run, like a sort order over a pruned column. + * Dropping it loses nothing today when Spark can call the transform's function, since no + * operator then requires an ordering over the transform. The write path sorts by the function + * call instead. Keeping such a sort order would add comparisons nobody uses. Dropping it also + * handles a transform Spark cannot evaluate. Revisit this if an ordering over a transform + * becomes a real requirement. */ override def outputOrdering: Seq[SortOrder] = { + def holdsTransform(e: Expression): Boolean = e.exists(_.isInstanceOf[TransformExpression]) (ordering, outputPartitioning) match { case (Some(o), p) => - val (prefix, rest) = o.span(_.references.subsetOf(outputSet)) + val (prefix, rest) = + o.span(order => order.references.subsetOf(outputSet) && !holdsTransform(order.child)) p match { case k: KeyedPartitioning if rest.nonEmpty => - val keyExprs = ExpressionSet(k.expressions) + val keyExprs = ExpressionSet(k.expressions.filterNot(holdsTransform)) prefix ++ rest.filter(order => keyExprs.contains(order.child)) case _ => prefix } case (_, k: KeyedPartitioning) if conf.v2BucketingPartitionKeyOrderingEnabled => - k.expressions.map(SortOrder(_, Ascending)) + k.expressions.filterNot(holdsTransform).map(SortOrder(_, Ascending)) case _ => Seq.empty } } diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala index 79973b599edec..84eac4627608c 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala @@ -24,7 +24,6 @@ import org.apache.spark.{Partition, SparkException} import org.apache.spark.rdd.{CoalescedRDD, PartitionCoalescer, PartitionGroup, RDD, SortedMergeCoalescedRDD} import org.apache.spark.sql.catalyst.InternalRow import org.apache.spark.sql.catalyst.expressions._ -import org.apache.spark.sql.catalyst.expressions.codegen.GenerateOrdering import org.apache.spark.sql.catalyst.plans.QueryPlan import org.apache.spark.sql.catalyst.plans.physical.{IdentityReducer, KeyedPartitioning, KeyReducer, Partitioning, UnknownPartitioning} import org.apache.spark.sql.catalyst.util.{truncatedString, InternalRowComparableWrapper} @@ -347,8 +346,8 @@ case class GroupPartitionsExec( } /** - * The ordering used by the k-way merge in [[SortedMergeCoalescedRDD]]. The generated comparator - * ([[GenerateOrdering]]) only needs each [[SortOrder]]'s sort key (child, direction, null + * The ordering used by the k-way merge in [[SortedMergeCoalescedRDD]]. The comparator that + * [[RowOrdering]] builds only needs each [[SortOrder]]'s sort key (child, direction, null * ordering), so `sameOrderExpressions` -- planner-only metadata that would otherwise be * serialized with the RDD in every task -- is dropped. */ @@ -360,7 +359,7 @@ case class GroupPartitionsExec( sparkContext.emptyRDD } else if (usesSortedMerge) { val partitionCoalescer = new GroupedPartitionCoalescer(groupedPartitions.map(_._2)) - val rowOrdering = new LazyCodeGenOrdering(kWayMergeOrdering, child.output) + val rowOrdering = new LazyRowOrdering(kWayMergeOrdering, child.output) new SortedMergeCoalescedRDD[InternalRow]( child.execute(), groupedPartitions.size, @@ -412,11 +411,9 @@ case class GroupPartitionsExec( outputPartitioning match { case p: Partitioning with Expression if reducers.isEmpty && conf.v2BucketingPreserveKeyOrderingOnCoalesceEnabled => - // Without reducers all merged partitions share the same original key value, so the key - // expressions remain constant within the output partition. The child's outputOrdering - // should already be in sync with the partitioning (either reported by the source or - // derived from it in DataSourceV2ScanExecBase), so we only need to keep the sort orders - // whose expression is a partition key expression -- all others are lost by concatenation. + // Without reducers all merged partitions share the same original key value, so the sort + // orders on key expressions still hold. The transform keys match nothing here, since + // `DataSourceV2ScanExecBase.outputOrdering` drops every sort order over a transform. val keyedPartitionings = p.collect { case k: KeyedPartitioning => k } val keyExprs = ExpressionSet(keyedPartitionings.flatMap(_.expressions)) child.outputOrdering.filter(order => keyExprs.contains(order.child)) @@ -489,15 +486,15 @@ class GroupedPartitionCoalescer( } /** - * A serializable [[Ordering]] for [[InternalRow]] that generates code-compiled comparison logic - * lazily on first use. The [[SortOrder]] expressions and output schema are serialized with the - * RDD; the generated comparator is rebuilt on the executor on first comparison via - * [[GenerateOrdering]]. + * A serializable [[Ordering]] for [[InternalRow]]. The [[SortOrder]] expressions and output schema + * are serialized with the RDD. The comparator is built on the executor, at the first comparison. + * [[RowOrdering]] builds it, as for `SortExec`. It tries generated code first and falls back to + * interpreted evaluation when that fails. */ -private class LazyCodeGenOrdering( +private class LazyRowOrdering( sortOrders: Seq[SortOrder], schema: Seq[Attribute]) extends Ordering[InternalRow] with Serializable { - @transient private lazy val generated: Ordering[InternalRow] = - GenerateOrdering.generate(sortOrders, schema) - override def compare(x: InternalRow, y: InternalRow): Int = generated.compare(x, y) + @transient private lazy val ordering: Ordering[InternalRow] = + RowOrdering.create(sortOrders, schema) + override def compare(x: InternalRow, y: InternalRow): Int = ordering.compare(x, y) } 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 6904476f754a7..7eb70aea4393b 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 @@ -41,6 +41,7 @@ import org.apache.spark.sql.execution.{ SortExec, SparkPlan, UnionExec} +import org.apache.spark.sql.execution.aggregate.SortAggregateExec import org.apache.spark.sql.execution.datasources.v2.{BatchScanExec, DataSourceV2ScanRelation, GroupPartitionsExec} import org.apache.spark.sql.execution.exchange.{ShuffleExchangeExec, ShuffleExchangeLike, ValidateRequirements} import org.apache.spark.sql.execution.joins.{BroadcastHashJoinExec, ShuffledHashJoinExec, SortMergeJoinExec} @@ -5511,6 +5512,119 @@ class KeyGroupedPartitioningSuite extends DistributionAndOrderingSuiteBase with } } + test("SPARK-59995: a scan's output ordering drops sort orders that hold a partition transform") { + // `t1` reports a transform in the middle, so the leading run of its ordering stops there. + // `t2` reports a transform key first. Each split holds a single key, so the sort order on `id` + // after it still holds. `t3` reports no ordering, so the scan derives one from its keys. In + // each case a sort on `id` right above the scan needs no `SortExec`. A `GROUP BY id` plans a + // `SortAggregateExec`. The query selects `ts`, so the transform's column stays in the scan + // output. Otherwise SPARK-59899's pruned-column handling would give the same orderings even + // without this fix. + val table1 = "transform_order_t1" + val table2 = "transform_order_t2" + val table3 = "transform_order_t3" + def asc(expr: Expression): SortOrder = sort(expr, SortDirection.ASCENDING) + createTable(table1, columns, Array(identity("id")), + Array(asc(FieldReference("id")), asc(years("ts")), asc(FieldReference("data")))) + createTable(table2, columns, Array(years("ts"), identity("id")), + Array(asc(years("ts")), asc(FieldReference("id")))) + createTable(table3, columns, Array(days("ts"), identity("id"))) + withSQLConf( + SQLConf.V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key -> "true", + SQLConf.REPLACE_HASH_WITH_SORT_AGG_ENABLED.key -> "true") { + Seq(table1, table2, table3).foreach { table => + sql(s"INSERT INTO testcat.ns.$table VALUES (1, 'aa', cast('2020-01-01' as timestamp))") + val df = sql(s"SELECT id, data, ts FROM testcat.ns.$table") + val scan = collectScans(df.queryExecution.executedPlan).head + assert(scan.output.exists(_.name == "ts"), s"test setup: ts stays in the output of $table") + assert(scan.outputOrdering.map(_.child.sql) === Seq("id"), table) + val sorted = df.sortWithinPartitions("id").queryExecution.executedPlan + assert(collect(sorted) { case s: SortExec => s }.isEmpty, table) + val aggregated = sql(s"SELECT id, max(ts) FROM testcat.ns.$table GROUP BY id") + .queryExecution.executedPlan + assert(collect(aggregated) { case a: SortAggregateExec => a }.nonEmpty, table) + } + } + } + + /** + * Inserts three items, two of them with `id` 1 on their own splits, so grouping the splits by + * `id` coalesces them. Then joins `items` with `purchases` on `joinCondition` and checks that + * the plan k-way merges over `expectedOrdering`, the columns the scan's ordering keeps before + * its transform. The query selects `arrive_time`, so the transform's column stays in the scan + * output. Otherwise SPARK-59899's pruned-column handling would drop the transform even without + * this fix. + */ + private def checkKWayMergeBeforeTransform( + joinCondition: String, + expectedOrdering: Seq[String]): Unit = { + sql(s"INSERT INTO testcat.ns.$items VALUES " + + "(1, 'aa', 10.0, cast('2021-01-01' as timestamp)), " + + "(1, 'ab', 11.0, cast('2022-01-01' as timestamp)), " + + "(2, 'bb', 20.0, cast('2021-01-01' as timestamp))") + val df = sql( + s""" + |${selectWithMergeJoinHint("i", "p")} + |i.id, i.name, i.arrive_time + |FROM testcat.ns.$items i JOIN testcat.ns.$purchases p ON $joinCondition + |""".stripMargin) + checkAnswer(df, Seq( + Row(1, "aa", Timestamp.valueOf("2021-01-01 00:00:00")), + Row(1, "ab", Timestamp.valueOf("2022-01-01 00:00:00")), + Row(2, "bb", Timestamp.valueOf("2021-01-01 00:00:00")))) + val merging = collectAllGroupPartitions(df.queryExecution.executedPlan) + .filter(_.enableSortedMerge) + assert(merging.length == 1, "expected one k-way merge") + assert(merging.head.child.output.exists(_.name == "arrive_time"), + "test setup: arrive_time stays in the scan output") + assert(merging.head.child.outputOrdering.map(_.child.sql) == expectedOrdering) + assert(merging.head.execute().isInstanceOf[SortedMergeCoalescedRDD[_]]) + } + + test("SPARK-59995: k-way merge over a reported ordering with a partition transform") { + // The join on (id, name) needs the merge, since the key ordering on id alone is not enough. + // The scan reports [id, name, years(arrive_time)] and keeps [id, name]. + val itemOrdering = Array( + sort(FieldReference("id"), SortDirection.ASCENDING), + sort(FieldReference("name"), SortDirection.ASCENDING), + sort(years("arrive_time"), SortDirection.ASCENDING)) + createTable(items, itemsColumns, Array(identity("id")), itemOrdering) + val namedPurchasesColumns = Array( + Column.create("item_id", LongType), + Column.create("name", StringType)) + createTable(purchases, namedPurchasesColumns, Array(identity("item_id"))) + sql(s"INSERT INTO testcat.ns.$purchases VALUES (1, 'aa'), (1, 'ab'), (2, 'bb')") + + withSQLConf( + SQLConf.REQUIRE_ALL_CLUSTER_KEYS_FOR_CO_PARTITION.key -> "false", + SQLConf.V2_BUCKETING_PRESERVE_ORDERING_ON_COALESCE_ENABLED.key -> "true") { + checkKWayMergeBeforeTransform("p.item_id = i.id AND p.name = i.name", Seq("id", "name")) + } + } + + test("SPARK-59995: k-way merge over an ordering derived from a partition transform key") { + // The scan reports no ordering, so it derives [id] from its keys, leaving out + // days(arrive_time). The merge must not call the transform's function. Spark cannot call this + // `days` at all, since `DaysFunction` implements neither `invoke` nor `produceResult`. The join + // on id projects the keys to id. With preserveKeyOrderingOnCoalesce off, the coalesced + // partitions keep no ordering on id unless they are merged. Before Spark 4.4, a join on a + // subset of the partition keys also needs requireAllClusterKeysForCoPartition off. + createTable(items, itemsColumns, Array(identity("id"), days("arrive_time"))) + createTable(purchases, purchasesColumns, Array(identity("item_id"))) + sql(s"INSERT INTO testcat.ns.$purchases VALUES " + + "(1, 10.0, cast('2021-01-01' as timestamp)), " + + "(2, 20.0, cast('2021-01-01' as timestamp))") + + withSQLConf( + SQLConf.REQUIRE_ALL_CLUSTER_KEYS_FOR_CO_PARTITION.key -> "false", + SQLConf.V2_BUCKETING_ALLOW_JOIN_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "true", + SQLConf.V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key -> "true", + SQLConf.V2_BUCKETING_PRESERVE_ORDERING_ON_COALESCE_ENABLED.key -> "true", + SQLConf.V2_BUCKETING_PRESERVE_KEY_ORDERING_ON_COALESCE_ENABLED.key -> "false") { + checkKWayMergeBeforeTransform("p.item_id = i.id", Seq("id")) + } + } + test("SPARK-56549: k-way merge enabled only when parent requires ordering") { // Dynamic gate: with the config enabled, k-way merge must be activated only when the parent // actually requires ordering (SMJ), and must stay off when the parent does not (hash join). diff --git a/sql/core/src/test/scala/org/apache/spark/sql/connector/WriteDistributionAndOrderingSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/connector/WriteDistributionAndOrderingSuite.scala index ec1b34b8c2101..2a2c0a7bb5564 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/connector/WriteDistributionAndOrderingSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/connector/WriteDistributionAndOrderingSuite.scala @@ -1545,10 +1545,12 @@ class WriteDistributionAndOrderingSuite extends DistributionAndOrderingSuiteBase val df = sql("SELECT id, data FROM testcat.ns1.test_table") val scans = collect(df.queryExecution.executedPlan) { case s: BatchScanExec => s } assert(scans.size === 1) - val ordering = scans.head.outputOrdering + val ordering = scans.head.ordering.getOrElse(Seq.empty) assert(ordering.nonEmpty, - "scan should report non-empty outputOrdering via SupportsReportOrdering") + "scan should keep the ordering reported via SupportsReportOrdering") assert(ordering.head.child.isInstanceOf[TransformExpression], "bucket-based sort order should resolve to a TransformExpression") + // SPARK-59995: the output ordering drops the sort order over the transform. + assert(scans.head.outputOrdering.isEmpty) } } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala index 62e316d632b95..95c9a45f505fa 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala @@ -17,16 +17,19 @@ package org.apache.spark.sql.execution.datasources.v2 +import org.apache.spark.{SparkConf, SparkException} import org.apache.spark.rdd.RDD +import org.apache.spark.serializer.JavaSerializer import org.apache.spark.sql.catalyst.InternalRow -import org.apache.spark.sql.catalyst.expressions.{Ascending, Attribute, AttributeReference, SortOrder, TransformExpression} +import org.apache.spark.sql.catalyst.expressions.{Ascending, Attribute, AttributeReference, CodegenObjectFactoryMode, SortOrder, TransformExpression} import org.apache.spark.sql.catalyst.plans.physical.{KeyedPartitioning, KeyReducer, Partitioning, PartitioningCollection, UnknownPartitioning} +import org.apache.spark.sql.catalyst.util.DateTimeTestUtils.date import org.apache.spark.sql.catalyst.util.InternalRowComparableWrapper -import org.apache.spark.sql.connector.catalog.functions.{BucketFunction, BucketReducer, Reducer} +import org.apache.spark.sql.connector.catalog.functions.{BucketFunction, BucketReducer, Reducer, YearsFunction} import org.apache.spark.sql.execution.{DummySparkPlan, LeafExecNode, SafeForKWayMerge} import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.test.SharedSparkSession -import org.apache.spark.sql.types.{DataType, IntegerType} +import org.apache.spark.sql.types.{DataType, IntegerType, TimestampType} class GroupPartitionsExecSuite extends SharedSparkSession { @@ -83,7 +86,7 @@ class GroupPartitionsExecSuite extends SharedSparkSession { test("SPARK-58324: k-way merge ordering drops sameOrderExpressions") { // The child ordering carries sameOrderExpressions (planner metadata). The k-way merge // comparator only needs the sort key, so kWayMergeOrdering keeps child/direction/nullOrdering - // but drops sameOrderExpressions, so LazyCodeGenOrdering does not serialize them with the RDD. + // but drops sameOrderExpressions, so LazyRowOrdering does not serialize them with the RDD. val childOrdering = Seq(SortOrder(exprA, Ascending, Seq(exprB, exprC))) val child = DummySparkPlan( outputPartitioning = KeyedPartitioning(Seq(exprA), Seq(row(1), row(2), row(1))), @@ -96,6 +99,34 @@ class GroupPartitionsExecSuite extends SharedSparkSession { assert(merged.forall(_.sameOrderExpressions.isEmpty)) } + test("SPARK-59995: the k-way merge ordering falls back to interpreted evaluation") { + // The scan drops every sort order over a transform, so none reaches the merge. Here a transform + // only stands in for a sort key without generated code. `TransformExpression` generates none, + // while `eval` calls its function. The ordering is serialized before any comparison, as the + // RDD ships it. + val ts = AttributeReference("ts", TimestampType)() + val ordering = new LazyRowOrdering( + Seq(SortOrder(TransformExpression(YearsFunction, Seq(ts)), Ascending)), Seq(ts)) + val serializer = new JavaSerializer(new SparkConf()).newInstance() + def ship(): LazyRowOrdering = + serializer.deserialize[LazyRowOrdering](serializer.serialize(ordering)) + // Under `CODEGEN_ONLY` the comparator cannot be built, which is the premise of this test. + val error = withSQLConf( + SQLConf.CODEGEN_FACTORY_MODE.key -> CodegenObjectFactoryMode.CODEGEN_ONLY.toString) { + intercept[SparkException] { + ship().compare(InternalRow(date(2022)), InternalRow(date(2021))) + } + } + assert(error.getMessage.contains("Cannot generate code for expression")) + // `FALLBACK` is the production default. + val shipped = ship() + withSQLConf(SQLConf.CODEGEN_FACTORY_MODE.key -> CodegenObjectFactoryMode.FALLBACK.toString) { + assert(shipped.compare(InternalRow(date(2022)), InternalRow(date(2021, 6))) > 0) + // Two timestamps in the same year compare equal, so the year is what is compared. + assert(shipped.compare(InternalRow(date(2021)), InternalRow(date(2021, 6))) === 0) + } + } + test("SPARK-56241: coalescing without reducers keeps key-expression orders from child") { // Key 1 appears on partitions 0 and 2, causing coalescing. val partitionKeys = Seq(row(1), row(2), row(1))