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]