andygrove commented on code in PR #5957:
URL: https://github.com/apache/datafusion-comet/pull/5957#discussion_r4109345999


##########
spark/src/main/scala/org/apache/comet/rules/RevertNativeForTransitionHeavyStages.scala:
##########
@@ -142,16 +151,25 @@ case class RevertNativeForTransitionHeavyStages(session: 
SparkSession)
   }
 
   /**
-   * Like `transformDown`, never descends stage-boundary children.
+   * Like `transformDown`, never descends stage-boundary children. If the rule 
rewrites the
+   * current node, re-apply it to the result so stacked transitions such as
+   * `CometSparkToColumnarExec(CometNativeColumnarToRowExec(x))` are fully 
unwrapped before
+   * children are visited. Spark's `transformDown` does not do this; leaving 
the inner C2R in
+   * place later calls `CometNativeColumnarToRowExec.withNewChildren` with a 
reverted row-based
+   * child, which asserts `child.supportsColumnar`.
    */
   private def transformStageDown(plan: SparkPlan)(
       rule: PartialFunction[SparkPlan, SparkPlan]): SparkPlan = {
     val transformed = rule.applyOrElse(plan, identity[SparkPlan])
-    val newChildren = transformed.children.map { child =>
-      if (isStageBoundary(child)) child else transformStageDown(child)(rule)
+    if (transformed ne plan) {
+      transformStageDown(transformed)(rule)
+    } else {
+      val newChildren = transformed.children.map { child =>
+        if (isStageBoundary(child)) child else transformStageDown(child)(rule)

Review Comment:
   With AQE off this still descends into the stage below an exchange. The 
boundary check only runs on the children of the transformed node. So when the 
stripped transition sits directly on a shuffle, the recursion strips the 
transitions inside the next stage down, and nothing puts them back, because 
`transformStageUp` and `insertTransitions` stop at the exchange. This is the 
problem in #6152, and here's a concrete repro: a copy-on-write `DELETE ... 
WHERE id IN (SELECT ...)` on a partitioned table, with 
`spark.sql.adaptive.enabled=false` and `maxTransitions=0`, fails the shuffle 
map stage with `ColumnarBatch cannot be cast to InternalRow`. Since this PR 
already rewrites `transformStageDown`, could it return `transformed` unchanged 
when it is a stage boundary, with `if (isStageBoundary(transformed)) 
transformed else transformStageDown(transformed)(rule)`? With that, the same 
DELETE commits once with the right rows and the same partition layout as the 
native write, and the rest of the s
 uite still passes.



##########
spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala:
##########
@@ -605,6 +605,37 @@ class CometIcebergWriteActionSuite
     }
   }
 
+  for (adaptive <- Seq(false, true)) {
+    test(s"transition-heavy fallback preserves Iceberg writes with 
AQE=$adaptive") {

Review Comment:
   The end-to-end case only covers an unpartitioned `INSERT ... VALUES`. There 
is no exchange under the write, and the restored `IcebergWriteExec` has no 
`ReplaceDataDispatchInfo`. Could we add a copy-on-write `DELETE` against a 
partitioned table, with AQE on and off, and compare its rows and partition 
directories with a sibling table written natively? That is the one shape where 
the restored node carries state the unit tests build by hand. With AQE off it's 
also the case that catches the boundary problem above.



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