dongjoon-hyun commented on code in PR #59251:
URL: https://github.com/apache/spark/pull/59251#discussion_r4209704835


##########
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:
   Nit: these assertions hold whether or not code generation failed, so this 
pins the fallback only while `TransformExpression.doGenCode` throws. If a 
transform generates code again, as in the round-1 version of this PR, the test 
would still pass without any fallback. And since the scan no longer lets a 
transform reach the merge, the transform here only stands in for a sort key 
without generated code, which the comment at L163-165 could say.
   
   Could we also check that a fresh deserialized copy fails under the session's 
`CODEGEN_ONLY`, e.g. `intercept[SparkException]` with "Cannot generate code for 
expression", so the test fails if the premise changes?



##########
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:
   Nit: the two-argument `sort(expr, SortDirection.ASCENDING)` already uses 
`NullOrdering.NULLS_FIRST`, the default null ordering of 
`SortDirection.ASCENDING`, and this suite already uses that form (e.g. L4043).
   
   Could this be `def asc(expr: Expression): SortOrder = sort(expr, 
SortDirection.ASCENDING)`?



##########
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:
   **`kWayMergeOrdering` is read from the child again at execution time, and 
the scan's derived part of it follows `partitionKeyOrdering.enabled` live, so a 
merge planned over a derived ordering can run with an empty one.** 
`DataSourceV2ScanExecBase.outputOrdering` derives the key ordering only while 
`spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` is on 
(DataSourceV2ScanExecBase.scala:168), and both this line and 
`kWayMergeIsFeasible` (L247-248) read `child.outputOrdering` when the node 
executes. For example, take a table identity-partitioned by `(a, b)` that 
reports no ordering, with `allowKeysSubsetOfPartitionKeys` and 
`preserveOrderingOnCoalesce` on, and plan `SELECT a, b, c, row_number() OVER 
(PARTITION BY a ORDER BY b) FROM t`: the plan has 
`GroupPartitionsExec(SortedMerge: true)` under the window and no `SortExec`. If 
the conf is turned off before `collect()`, the node can still merge, since its 
`usesSortedMerge` may already have been evaluated during planning, bu
 t `kWayMergeOrdering` is now empty, so `LazyRowOrdering(Nil)` treats every row 
as equal and the window gets its rows out of order. This predates the PR 
(SPARK-56241) and needs two non-default confs plus a conf change after 
planning, so it is the same kind of issue as SPARK-59279 rather than a blocker 
here.
   
   Could we file a follow-up JIRA to freeze the merge ordering when 
`tryEnableSortedMerge` decides on the merge, as SPARK-59279 did for the merge 
config?



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