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


##########
sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala:
##########
@@ -6571,6 +6571,102 @@ class KeyGroupedPartitioningSuite
     }
   }
 
+  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`.
+    val table1 = "transform_order_t1"
+    val table2 = "transform_order_t2"
+    val table3 = "transform_order_t3"
+    def asc(expr: Expression): SortOrder =
+      sort(expr, SortDirection.ASCENDING, NullOrdering.NULLS_FIRST)
+    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")))
+    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")

Review Comment:
   Added the setup asserts in 6eca5ef5e41, and said in the scaladoc of 
`checkKWayMergeBeforeTransform` why the column has to stay.
   



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala:
##########
@@ -335,7 +334,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)

Review Comment:
   Agreed. I filed SPARK-60044 for it.
   



##########
sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala:
##########
@@ -157,6 +159,22 @@ class GroupPartitionsExecSuite extends SharedSparkSession {
     assert(merged.forall(_.sameOrderExpressions.isEmpty))
   }
 
+  test("SPARK-59995: the k-way merge ordering falls back to interpreted 
evaluation") {
+    // A `TransformExpression` generates no code, while `eval` calls its 
function. The ordering is
+    // serialized before any comparison, as the RDD ships it. Test sessions 
run `CODEGEN_ONLY`, so
+    // this sets the production default, `FALLBACK`.
+    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()
+    val shipped = 
serializer.deserialize[LazyRowOrdering](serializer.serialize(ordering))
+    withSQLConf(SQLConf.CODEGEN_FACTORY_MODE.key -> 
CodegenObjectFactoryMode.FALLBACK.toString) {

Review Comment:
   Added in 6eca5ef5e41. A fresh copy now fails under the session's 
`CODEGEN_ONLY` with "Cannot generate code for expression", and the comment says 
that the transform only stands in for a sort key without generated code.
   



##########
sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala:
##########
@@ -6571,6 +6571,102 @@ class KeyGroupedPartitioningSuite
     }
   }
 
+  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`.
+    val table1 = "transform_order_t1"
+    val table2 = "transform_order_t2"
+    val table3 = "transform_order_t3"
+    def asc(expr: Expression): SortOrder =
+      sort(expr, SortDirection.ASCENDING, NullOrdering.NULLS_FIRST)

Review Comment:
   Done in 6eca5ef5e41, and in the reported merge test's sort orders too.
   



-- 
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