Conversation
Spark's InMemoryTableScanExec lists its InMemoryRelation, and through it the cached plan, as inner children, so EXPLAIN draws the plan that built the cache below the scan. Do the same for CometInMemoryTableScanExec. ExtendedExplainInfo walks inner children, so leave the scan's out of Comet's coverage and fallback reporting, as it already does for CometEmptyRelationExec. The cached plan runs when the relation is materialized, not as part of every query that reads it. Closes apache#6572.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Comet cache scans omitted the cached plan from EXPLAIN, the SQL graph and event-log plans.
- Design approach: Expose the relation through
innerChildrenand the cached physical plan throughsubqueries, while excluding both from Comet’s execution reporting. - Correctness / compatibility analysis: Checked Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. The graph construction works, but the generic
subqueriesoverride introduces the Spark 4.2 metric regression below. - Key design decisions: Executable children remain unchanged, and the lazy override accommodates Spark 3.x. The implementation is small, but using an execution-related traversal API for display data has observable consequences beyond the UI. Additional traversal is identifiable, but no performance regression was measured.
- Implementation sketch: Two scan overrides, one reporting exclusion, a test helper and AQE-on/off regression tests.
- Behavioral changes worth calling out: Compared with
branch-1.1, displaying the cached subtree is intended. Codegen inspection also reaches that subtree, and Spark 3.x can emit another final plan update. The last-attempt metric failure is unintended. - Suggested improvements: Address the P2 finding and cover a metric recorded outside an already materialized cache containing a Comet shuffle.
Reviewed the entire four-file diff from f980d6fb59f51ae20c0a6916eec9771ace8e50d1 to a422a549d327933010fc04c59bc29ec68bef2863, including prerequisite commit 81736558d43e00f55499192761c6cd9303c31444. Verified its reusable source evidence. The PR remains non-draft. Snapshot and live discussion checks contained no existing review concerns. Routed skill: .ai/skills/review-comet-pr/SKILL.md. No sibling area skill applies.
Exact-head CI at 2026-10-03 19:19 UTC: 38 successful, 7 running and 35 skipped checks, with no failures. Spark 4.1 and 3.4 execution jobs passed. The Spark 4.1 log confirms both new tests passed within 1,229 successful tests. Execution jobs for 3.5, 4.0 and 4.2 remained pending. Upstream Spark SQL suites were skipped.
Validation: git diff --check passed. A Spark 3.5.9 API probe passed eight cases covering AQE, nested caches and materialization. A bounded replay of Spark 4.2’s scope-extraction method confirmed the finding. These probes did not execute Comet native code. The Spark 4.1 probe compiled but could not start because the available runtime lacked KVStore. No local full Comet build ran. Build artifacts and dependencies were absent, and Maven Central access was blocked.
Review state: Request changes.
| // 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) |
There was a problem hiding this 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.
… and event log Spark builds the SQL tab's graph and the plans in the event log with SparkPlanInfo.fromSparkPlan. It gives its own InMemoryTableScanExec the cached plan as a child, but recognizes that scan by its class, and for any other node takes the children and the subqueries. Expose the cached plan as the one subquery of CometInMemoryTableScanExec. Spark runs subqueries from a plan's expressions and only walks this list, so the cached plan does not become part of the query that reads the cache. Closes apache#6463.
a422a54 to
5d5b520
Compare
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Comet cache scans omitted the cached subtree from EXPLAIN, the SQL graph and event-log plans.
- Design approach: Expose the relation through
innerChildrenand its cached physical plan throughsubqueries, while excluding it from Comet’s coverage and fallback reporting. - Correctness / compatibility analysis: Checked Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. The display behavior follows Spark’s implementation. The existing Spark 4.2 metric regression remains unresolved: a metric recorded outside a materialized cache can return
Noneinstead ofSome(100)because traversal now encounters a cached Comet shuffle. - Key design decisions: Executable children remain unchanged, and the lazy override supports both Spark 3.x and 4.x. The implementation is small, but using
subqueriesfor display couples it to other plan walkers. No additional reproducible performance regression was identified. - Implementation sketch: Two scan overrides, one reporting exclusion, a test helper and AQE-on/off regression tests.
- Behavioral changes worth calling out: Compared with
branch-1.1, displaying the cached subtree is intended. Codegen inspection also reaches it, and Spark 3.x can emit an additional final plan update. The last-attempt metric regression is unintended. - Suggested improvements: Address the existing metric-traversal concern by isolating the display subtree or providing compatible handling, with a regression covering a metric outside an already materialized cache.
Reviewed full SHA 5d5b52045d069dae6e23be1ab2106f572f67680d, including all four files and both prerequisite commits from #6574. The requested base is 569eaa59d032f758669777964aee2d24eb55ebae; the PR merge-base is f980d6fb59f51ae20c0a6916eec9771ace8e50d1. Read AGENTS.md and routed through .ai/skills/review-comet-pr/SKILL.md. No sibling area skill applies. The PR remains non-draft.
Read the snapshot and live discussions. No additional introduced P1/P2 issues found within this review. The existing P2 concern remains at spark/src/main/scala/org/apache/spark/sql/comet/CometInMemoryTableScanExec.scala:95. Production behavior is unchanged from the previously reviewed head, so this is not duplicated as a new finding.
Validation: git diff --check passed. Reran the Spark 3.5.9 API probe successfully across eight AQE, nested-cache and materialization cases. Verified the Spark 4.2 scope-extraction replay byte-for-byte against upstream source and reproduced the change from Right(List(3)) to Left(Unsupported ShuffleExchangeLike: ...). These probes validate Spark traversal behavior, not end-to-end Comet execution. No local Comet/native build or suite ran because this checkout lacks built artifacts and Maven Spark dependencies.
Exact-head CI at 2026-10-03 20:45 UTC: 6 successful, 2 running and 14 skipped checks, with no failures reported. Linux lint remained running, and no exact-head Comet test verdict was available. run-all-spark-profiles is applied. Upstream Spark SQL suites were skipped.
Review state: Request changes remains warranted for the existing P2 concern.
Which issue does this PR close?
Closes #6463.
Rationale for this change
Spark builds the SQL tab's graph, and the plans the event log records, with
SparkPlanInfo.fromSparkPlan. It gives its ownInMemoryTableScanExecthe cached plan as a child, but recognizes that scan by its class. For any other node it takesplan.children ++ plan.subqueries, so for a relation cached in Comet's format the tree ended atCometInMemoryTableScan, and the cached plan's metrics were missing below it.The issue expected that Comet could not fix this alone, because a child would make the cached plan part of the query that reads the cache. A subquery does not: Spark runs subqueries from a plan's expressions and only walks the
subquerieslist.What changes are included in this PR?
This PR is stacked on #6574. The first two commits are that PR's, so review the last two. It needs #6574's explicit
innerChildrenoverride:QueryPlan.innerChildrendefaults tosubqueries, so without it the cached plan would also land in EXPLAIN without itsInMemoryRelationline, and in the coverage and fallback reporting ofExtendedExplainInfo.CometInMemoryTableScanExecoverridessubqueriesto return the cached plan. It is alazy valbecause Spark 3.x declaressubqueriesas one, and alazy valoverrides Spark 4'sdefas well.CometSparkPlanInfoHelper, because theSparkPlanInfocompanion object isprivate[execution].Other readers of a physical plan's
subqueriesnow follow it into the cached plan too, which they do not do for Spark's own scan. Checked against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0:CollectMetricsExec.collect, throughAdaptiveSparkPlanHelper.collectWithSubqueries, would find observed metrics in the cached plan, as it does below Spark's scan by matching its class. Nothing changes today, because fix: keep Spark's cache scan for a relation whose cached plan records observed metrics #6421 keeps Spark's scan for any relation whose cached plan records observed metrics.debugCodegenoutput includes the cached plan's codegen stages.AdaptiveSparkPlanExec.finalPlanUpdateposts one more plan update when the final plan contains the scan outside a query stage.SQLLastAttemptAccumulatorwalks the cached plan too. Its documentation already declares cached plans undefined behavior, and it bails out on Comet's shuffle exchanges anyway.CacheManager.validateCachedEntryForTransactioncould register one more scan for a nested cache in a DSv2 transaction.How are these changes tested?
A new test in
CometInMemoryCacheSuite, with AQE off and on, compares theSparkPlanInfosubtree below the scan with the cached plan's ownSparkPlanInfo. Without thesubqueriesoverride it fails, with no children below the scan.CometInMemoryCacheSuitepasses locally on the default Spark 4.1 profile (61 tests). Other profiles were not run locally, sorun-all-spark-profilesis applied for CI to compile and run the suites on them.Spark's
CachedTableSuitetest "SPARK-35332: Make cache plan disable configs configurable - check AQE", which #5634 skips under Comet, reads the cached plan from this tree and should now find it. Whether it passes was not checked.