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]