peter-toth commented on code in PR #59251:
URL: https://github.com/apache/spark/pull/59251#discussion_r4205062497


##########
sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala:
##########
@@ -6571,6 +6571,78 @@ class KeyGroupedPartitioningSuite
     }
   }
 
+  /** Two rows with `id` 1, each on its own split, so a join on `id` coalesces 
them. */
+  private def insertItemsInTwoYears(): 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))")
+
+  /**
+   * Joins `items` with `purchases` on `joinCondition` and checks that the 
plan k-way merges over an
+   * ordering with a partition transform.
+   */
+  private def checkKWayMergeOverTransform(joinCondition: String): Unit = {
+    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(

Review Comment:
   After 913f317a345 the merge no longer compares rows by the transform, so 
there is no transform order left to check. The tests now check the ordering the 
merge runs over, `[id, name]` and `[id]`.
   



##########
sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala:
##########
@@ -6571,6 +6571,78 @@ class KeyGroupedPartitioningSuite
     }
   }
 
+  /** Two rows with `id` 1, each on its own split, so a join on `id` coalesces 
them. */
+  private def insertItemsInTwoYears(): 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))")
+
+  /**
+   * Joins `items` with `purchases` on `joinCondition` and checks that the 
plan k-way merges over an
+   * ordering with a partition transform.
+   */
+  private def checkKWayMergeOverTransform(joinCondition: String): Unit = {
+    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.outputOrdering.exists(_.child.isInstanceOf[TransformExpression]),

Review Comment:
   Thanks, this is the simpler fix. I did it in the scan in 913f317a345, so 
every consumer of its ordering sees it, not only `GroupPartitionsExec`:
   - in a reported ordering, a transform ends the leading run, like a pruned 
column;
   - a transform key is dropped, while a sort order on another key still holds 
after it, since each partition holds a single key;
   - the derived ordering leaves the transform keys out.
   



##########
sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala:
##########
@@ -6571,6 +6571,78 @@ class KeyGroupedPartitioningSuite
     }
   }
 
+  /** Two rows with `id` 1, each on its own split, so a join on `id` coalesces 
them. */
+  private def insertItemsInTwoYears(): 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))")
+
+  /**
+   * Joins `items` with `purchases` on `joinCondition` and checks that the 
plan k-way merges over an
+   * ordering with a partition transform.
+   */
+  private def checkKWayMergeOverTransform(joinCondition: String): Unit = {
+    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.outputOrdering.exists(_.child.isInstanceOf[TransformExpression]),
+      "expected a partition transform in the merge's ordering")
+    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 merge's ordering is [id, name, years(arrive_time)], so generating 
its comparator
+    // generates code for the transform.
+    val itemOrdering = Array(
+      sort(FieldReference("id"), SortDirection.ASCENDING, 
NullOrdering.NULLS_FIRST),
+      sort(FieldReference("name"), SortDirection.ASCENDING, 
NullOrdering.NULLS_FIRST),
+      sort(years("arrive_time"), SortDirection.ASCENDING, 
NullOrdering.NULLS_FIRST))
+    createTable(items, itemsColumns, Array(identity("id")), itemOrdering)
+    insertItemsInTwoYears()
+    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") {
+      checkKWayMergeOverTransform("p.item_id = i.id AND p.name = i.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, years(arrive_time)] 
from its keys. The join
+    // on id projects the keys to id. With preserveKeyOrderingOnCoalesce off, 
the coalesced
+    // partitions keep no ordering on id unless they are merged.
+    createTable(items, itemsColumns, Array(identity("id"), 
years("arrive_time")))

Review Comment:
   After 913f317a345 the merge never calls a transform's function, so the way 
Spark calls it does not matter there. The derived test uses `days`, which Spark 
cannot call at all.
   



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to