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]