andygrove commented on code in PR #5957:
URL: https://github.com/apache/datafusion-comet/pull/5957#discussion_r4184571124
##########
spark/src/main/scala/org/apache/comet/rules/RevertNativeForTransitionHeavyStages.scala:
##########
@@ -62,32 +63,41 @@ case class RevertNativeForTransitionHeavyStages(session:
SparkSession, wholePlan
plan match {
case _: BroadcastExchangeLike => plan
case exchange: ShuffleExchangeLike =>
- revertStageIfNeeded(exchange.child, exchange.supportsColumnar)
+ revertShuffleStageIfNeeded(exchange)
.map(reverted => exchange.withNewChildren(Seq(reverted)))
.getOrElse(plan)
case _ =>
- // Result stage: its output is collected as rows, so no consumer
requires columnar input
- // and the reverted stage needs no trailing R2C.
+ // Result stage: its output is collected as rows.
revertStageIfNeeded(plan, outputColumnar = false).getOrElse(plan)
Review Comment:
The result stage still assumes it has to produce rows, and that isn't true
for a cached query on Spark 4.0+. When the cache serializer accepts columnar
input, which Comet's `ArrowCachedBatchSerializer` does, Spark marks the cached
`AdaptiveSparkPlanExec` columnar and AQE plans the final stage with
`outputsColumnar = true`. With the Comet cache on, AQE on,
`transitionRevert.enabled=true`, `maxTransitions=0` and
`spark.comet.exec.filter.enabled=false`, caching `SELECT _2, s FROM (SELECT _2,
sum(_1) AS s FROM tbl GROUP BY _2) WHERE s > 10` and collecting it fails with
`FilterExec has column support mismatch`. `main` fails the same way, so it
isn't a regression, but it's the same output-format restoration this PR is
fixing. Could both result-stage call sites (here and line 81) pass
`outputColumnar = plan.supportsColumnar && !plan.supportsRowBased` instead of
`false`? The root already says which format the consumer wants. I tried that
locally and the repro passes, along with all of `Re
vertNativeForTransitionHeavyStagesSuite`. The test would need a suite that
installs the cache serializer the way `CometInMemoryCacheSuite` does. If you'd
rather keep this PR to the write fix, could we file an issue for it instead?
##########
docs/source/contributor-guide/adding_a_new_operator.md:
##########
@@ -722,6 +722,27 @@ Use `QueryPlanSerde.exprToProto` to convert Spark
expressions to protobuf:
val protoExpr = exprToProto(sparkExpr, inputSchema)
```
+### Restoring the Spark operator (`sparkFallback`)
+
+`CometExec.originalPlan` is the Spark operator this node replaced.
`CometExecRule` copies
+`originalPlan.logicalLink` onto the Comet node, which is how AQE finds the
node again when it
+re-plans a stage. `RevertNativeForTransitionHeavyStages` calls
`sparkFallback(newChildren)` to
Review Comment:
This says the rule rebuilds the Spark operator through `sparkFallback`, but
`revertToSpark` handles `CometLocalTopKExec`, `CometNativeScanExec` and
`CometIcebergNativeScanExec` itself before it gets there
(`RevertNativeForTransitionHeavyStages.scala:263-281`). So `sparkFallback`
isn't the whole contract yet, and on `CometLocalTopKExec` the default would
return a second `TakeOrderedAndProjectExec` that applies the TopK twice. Could
those three become `sparkFallback` overrides, with the local TopK returning its
child and the two scans carrying their live DPP filters across? Then the rule
only ever calls `sparkFallback`, and this section holds for the next operator
whose live state differs from its `originalPlan`. If you'd rather not move the
code, could the section name the exceptions?
--
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]