cloud-fan commented on code in PR #58419:
URL: https://github.com/apache/spark/pull/58419#discussion_r4000506919
##########
sql/core/src/test/scala/org/apache/spark/sql/execution/UnionCodegenSuite.scala:
##########
@@ -628,6 +659,130 @@ class UnionCodegenSuite extends SharedSparkSession {
}
}
+ test("SPARK-59122: a fused union keeps numOutputRows and reports
UnknownPartitioning") {
+ // The children's partitioning is not stable while the plan is being
prepared:
+ // `InMemoryTableScanExec.cachedPlan` unwraps the inner
`AdaptiveSparkPlanExec` only once
+ // `isFinalPlan` is true, and reports `UnknownPartitioning(0)` until then,
so the union looks
+ // plain and is fused. The projection is what makes that reachable:
`supportsColumnar` is
+ // `children.forall`, so one row-based `ProjectExec` over the columnar
scan is enough to make
+ // it false, and without one `supportCodegenFailureReason` reports
`columnar` and nothing
+ // fuses. `SELECT *` or a plain alias collapses the projection away and
does not reproduce
+ // this. Once the cache stages finalise, both children report the same
`HashPartitioning`, and
+ // re-deriving the decision at that point left `metrics` empty while the
generated code still
+ // incremented it, so `doProduce` threw `key not found: numOutputRows`.
+ //
+ // Both halves of the decision are asserted here. Registering the metric
unconditionally would
+ // fix the crash and leave the other half broken: a fused union
concatenates its children's
+ // partitions, so claiming their `HashPartitioning` would let a parent
satisfy a clustered
+ // distribution from an RDD that does not have it.
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "true",
+ SQLConf.UNION_OUTPUT_PARTITIONING.key -> "true") {
+ withTempView("v") {
+ cacheAggregateView("v")
+ val df = spark.sql("SELECT k, abs(s) AS s FROM v UNION ALL SELECT k, s
FROM v")
+ // Execute this DataFrame rather than a count over it: the plan being
inspected has to be
+ // the one that ran, and an AQE plan that never ran has no final plan
to inspect.
+ assert(df.collect().length == 20)
+ val fused = fusedUnions(df)
+ assert(fused.nonEmpty,
+ "this shape must actually fuse, or the test is not exercising the
defect")
+ fused.foreach { u =>
+ assert(u.metrics.contains("numOutputRows"),
+ "a fused union must register the metric its generated code
increments")
+ assert(u.outputPartitioning.isInstanceOf[UnknownPartitioning],
+ s"a fused union must not claim a concrete partitioning, got
${u.outputPartitioning}")
+ }
+ }
+ }
+ }
+
+ test("SPARK-59122: a partitioning-aware union keeps its layout when the conf
changes between " +
+ "planning and execution") {
+ // `spark.sql.unionOutputPartitioning` is read where the plain-union
decision is latched, not on
+ // every `outputPartitioning` call, so a plan executes by the partitioning
it was planned
+ // against. Reading it per call let the parent aggregate lose its exchange
at planning and get a
+ // plain concatenation at execution, reporting each group twice. The
`checkAnswer` below stays
+ // outside the block that planned the DataFrame on purpose: the plan is
forced inside that
+ // block and `executedPlan` is memoized, so the two phases see different
confs. Asserting
+ // inside it, or dropping the second `withSQLConf`, makes the test pass
without testing this.
+ withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false") {
+ val left = spark.range(0, 20, 1, 2).selectExpr("id % 5 AS k")
+ val right = spark.range(20, 40, 1, 2).selectExpr("id % 5 AS k")
+
+ val planned = withSQLConf(SQLConf.UNION_OUTPUT_PARTITIONING.key ->
"true") {
+ val df = left.repartition(4, col("k"))
+ .union(right.repartition(4, col("k"))).groupBy("k").count()
+ val plan = df.queryExecution.executedPlan
+ val unions = plan.collect { case u: UnionExec => u }
+ assert(unions.size == 1)
+ // Not asserted through `isPlainUnion`: that call latches the
decision, which would warm
+ // a field-based implementation's memo and hide the regression this
test is for. The
+ // exchange count below proves the union reported a concrete
partitioning, without
+ // touching the node.
+ assert(plan.collect { case s: ShuffleExchangeExec => s }.size == 2,
+ "only the two repartitions may shuffle; the aggregate's exchange
must have been elided")
+ df
+ }
+
+ withSQLConf(SQLConf.UNION_OUTPUT_PARTITIONING.key -> "false") {
+ // Each side contributes four ids per `k`, so the answer is fixed.
Comparing against the
+ // same query run with the conf off would also pass if both paths
regressed to ten rows.
+ checkAnswer(planned, (0L until 5L).map(k => Row(k, 8L)))
+ }
+ }
+ }
+
+ test("SPARK-59122: a fused union keeps numOutputRows when the codegen conf
changes between " +
+ "planning and execution") {
+ // `supportCodegenFailureReason` used to read
`WHOLESTAGE_UNION_CODEGEN_ENABLED` live, and the
+ // copy that `insertInputAdapter` puts inside the codegen shell evaluated
it for the first time
+ // at execution. Planned with the conf on the union is fused, so the
generated code increments
+ // `numOutputRows`; if the copy re-derives the reason with the conf off,
`metrics` comes back
+ // empty and `doProduce` throws `key not found: numOutputRows`.
+ withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false") {
+ val planned = withSQLConf(SQLConf.WHOLESTAGE_UNION_CODEGEN_ENABLED.key
-> "true") {
+ // Each child is an exchange, which is not `CodegenSupport`, so
`insertInputAdapter` wraps
+ // it and `withNewChildren` really does produce a copy. With
codegen-support children it
Review Comment:
Confirmed: the comment now says CodegenSupport children can also lead to a
copy and uses the Exchange children only to guarantee a real InputAdapter
rebuild for this test. Resolved.
<!-- SPARK_DEV_REVIEW_REPLY
{"feedback_id":"inline:3993643925","thread_id":"inline:3993643925","verdict_sha256":"3f145606514841e98aa4f828110c8e9a2ef2a823fe82b8c1858abbafd71ce3c8"}
-->
--
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]