sunchao commented on code in PR #6577: URL: https://github.com/apache/datafusion-comet/pull/6577#discussion_r4174502176
########## spark/src/main/scala/org/apache/spark/sql/comet/CometInMemoryTableScanExec.scala: ########## @@ -80,6 +81,21 @@ case class CometInMemoryTableScanExec( // newlines, which breaks the tree of every plan that reads the cache. override def stringArgs: Iterator[Any] = Iterator(originalPlan) + // Spark's own scan lists its InMemoryRelation as an inner child, and the relation lists the + // cached plan, so EXPLAIN draws the plan that built the cache below the scan. Do the same. + // ExtendedExplainInfo leaves them out of Comet's own reporting: the cached plan runs when the + // relation is materialized, not as part of every query that reads it. + override def innerChildren: Seq[QueryPlan[_]] = Seq(originalPlan.relation) + + // SparkPlanInfo, which the SQL tab's graph and the event log's plans are built from, gives + // Spark's own scan its cached plan as a child, but recognizes that scan by its class. For any + // other node it takes the children and the subqueries, so expose the cached plan as the one + // subquery. A child would make the cached plan part of the query that reads the cache, but + // Spark only walks this list: subqueries run from a plan's expressions. Its other walkers, such + // as collectWithSubqueries, follow it into the cached plan too. A lazy val, because Spark 3.x + // declares subqueries as one, and a lazy val overrides Spark 4's def as well. + @transient override lazy val subqueries: Seq[SparkPlan] = Seq(originalPlan.relation.cachedPlan) Review Comment: [P2] Keep cached shuffles out of last-attempt metric traversal. On Spark 4.2 with Comet’s cache enabled, materialize `spark.range(0, 100, 1, 2).repartition(2).cache()`, then increment a `SQLLastAttemptMetrics.createMetric` accumulator in a subsequent `.map` over that cache. The consuming query has no shuffle, and the metric is entirely outside the cache, so `lastAttemptValueForDataset` should return `Some(100)`. This override makes `SQLLastAttemptAccumulator.extractStageRDDScopes` enter the cached plan, encounter `CometShuffleExchangeExec`, and return `Left(Unsupported ShuffleExchangeLike: ...)`. The metric accessor consequently returns `None`. Previously the cached shuffle was outside this traversal. Spark’s documented undefined behavior applies when the metric itself was used inside the cached plan, which does not cover this case. Please isolate display-only cached plans from this walker, or provide compatible handling, and add a regression for an outside-cache metric. Evidence: A bounded JVM probe replayed Spark v4.2.0’s `extractStageRDDScopes` method byte-for-byte using Spark 4.1.3 plan classes, an equivalent helper companion and stable substitute scope IDs. Its cached subtree contained a plugin exchange implementing `ShuffleExchangeLike`, matching Comet’s inheritance, beneath a consuming stage. Changing only whether the cache exposed that subtree through `subqueries` changed the result from `Right(List(3))` to `Left(Unsupported ShuffleExchangeLike: org.apache.spark.sql.execution.metric.PluginShuffle)`. Probe: `/tmp/pr6577-review-a422a549/ScopeRegressionProbe.scala`, output: `/tmp/pr6577-review-a422a549/scope-probe.log`. Spark v4.2.0 `SQLLastAttemptAccumulator.scala` lines 343–350 reject non-Spark shuffle implementations, lines 424–426 traverse these subqueries, and lines 257–263 convert the failure to `None`. This validates the traversal regression, not an end-to-end Comet execution. -- 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]
