Skip to content

fix: draw the cached plan below CometInMemoryTableScan in the SQL tab and event log - #6577

Merged
comphead merged 3 commits into
apache:mainfrom
comphead:cached-plan-sql-tab-6463
Oct 4, 2026
Merged

comphead merged 3 commits into
apache:mainfrom
comphead:cached-plan-sql-tab-6463

Conversation

@comphead

@comphead comphead commented Oct 3, 2026 •

Copy link
Copy Markdown
Contributor

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 draws its own InMemoryTableScanExec with the cached plan below it, but recognizes that scan by its class. For any other node it takes plan.children ++ plan.subqueries, so for a relation cached in Comet's format the tree ended at CometInMemoryTableScan, 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 subqueries list.

What changes are included in this PR?

It builds on #6574, now merged, and needs its explicit innerChildren override: QueryPlan.innerChildren defaults to subqueries, so without it EXPLAIN would draw this subquery as well, and ExtendedExplainInfo would count it.

  • CometInMemoryTableScanExec overrides subqueries to return the Spark InMemoryTableScanExec it replaces, not the cached plan. SparkPlanInfo draws the cached plan below that scan through its own special case, so the tree reads CometInMemoryTableScan, then Spark's scan of the cache, then the cached plan. It is a lazy val because Spark 3.x declares subqueries as one, and a lazy val overrides Spark 4's def as well.
  • A test-only CometSparkPlanInfoHelper, because the SparkPlanInfo companion object is private[execution].

Exposing Spark's scan rather than the cached plan leaves every other reader of subqueries where it is for Spark's own scan, because they treat that scan as a leaf or special-case it. An earlier revision exposed the cached plan, and Spark 4.2's SQLLastAttemptAccumulator then reached its Comet shuffle and gave up, so a last-attempt metric used outside the cache returned None (see the review thread). Checked against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0, two differences remain:

  • CollectMetricsExec.collect now reaches Spark's scan through the subquery and, as it does for that scan in Spark's own plans, collects observed metrics from its cached plan. 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.
  • On Spark 3.4 and 3.5, AdaptiveSparkPlanExec.finalPlanUpdate posts one more plan update when the final plan contains the scan outside a query stage, because it checks for any node with subqueries.

How are these changes tested?

A new test in CometInMemoryCacheSuite, with AQE off and on, checks that the scan's SparkPlanInfo node has Spark's scan below it, with the cached plan's own SparkPlanInfo below that. It also checks that collectWithSubqueries does not find the cached plan's shuffle in a query that has none of its own. The first check fails without the override, and the second fails when the cached plan itself is exposed, as in the earlier revision.

CometInMemoryCacheLastAttemptMetricSuite, under spark/src/test/spark-4.2, runs the scenario from the review: a last-attempt metric in a map over a cache whose plan has a Comet shuffle must still report Some(100). It is registered in both workflow files. It was not run locally, because the 4.2 profile cannot be built offline here, so run-all-spark-profiles gives its first result.

CometInMemoryCacheSuite passes locally on the default Spark 4.1 profile (61 tests). Spark's CachedTableSuite test "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.

@github-actions github-actions Bot added bug Something isn't working area:scan Parquet scan / data reading labels Oct 3, 2026
@comphead comphead added the run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue label Oct 3, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 innerChildren and the cached physical plan through subqueries, 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 subqueries override 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)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

@comphead comphead Oct 4, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks, confirmed and fixed in c190f6e (rebased onto main after #6574 merged). subqueries now returns Spark's own InMemoryTableScanExec (originalPlan) rather than the cached plan. SparkPlanInfo special-cases that class, so the SQL tab and event log still draw the cached plan, one level down, while SQLLastAttemptAccumulator and the other subquery walkers stop at that scan, as they do in Spark's own plans.

Your outside-cache scenario is now a 4.2-only CometInMemoryCacheLastAttemptMetricSuite: it checks that the cached plan has a Comet shuffle, then expects Some(100). A check on every version also asserts that collectWithSubqueries no longer reaches the cached plan's shuffle. The 4.2 suite could not run locally, so the run-all-spark-profiles run gives its first result.

@comphead
comphead force-pushed the cached-plan-sql-tab-6463 branch from a422a54 to 5d5b520 Compare October 3, 2026 20:41

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 innerChildren and its cached physical plan through subqueries, 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 None instead of Some(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 subqueries for 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.

… 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.
Spark 4.2's SQLLastAttemptAccumulator walks a plan's subqueries to find
a metric's stages and gives up on a shuffle it does not know. With the
cached plan as the subquery it reached the cached plan's Comet shuffle,
so a last-attempt metric used outside the cache returned None.

SparkPlanInfo draws the cached plan below Spark's own
InMemoryTableScanExec, matched by class, so expose that scan instead.
Every other walker of subqueries stops at it, as in Spark's own plans.

Adds a check on every version that collectWithSubqueries does not reach
the cached plan's shuffle, and a Spark 4.2 suite for the outside-cache
metric.
@comphead
comphead force-pushed the cached-plan-sql-tab-6463 branch from 1a36ec9 to c190f6e Compare October 4, 2026 16:44

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Comet cache scans omitted the cached plan and its metrics from the SQL graph and event-log plans.
  • Design approach: Expose the original Spark cache scan through subqueries. Spark’s existing SparkPlanInfo handling then displays the cached plan underneath it.
  • Correctness / compatibility analysis: Checked Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Preparation and transformations follow expressions, so the added scan does not execute. The current override addresses the existing Spark 4.2 metric concern by keeping the cached Comet shuffle outside scope extraction.
  • Key design decisions: A single lazy override supports Spark 3.x and 4.x without another production abstraction. The existing innerChildren override preserves EXPLAIN and Comet reporting. Graph construction and observed-metric collection gain traversal work, but no reproducible P1/P2 performance regression was identified.
  • Implementation sketch: One production override, a test helper, AQE-on/off graph and traversal assertions, and a Spark 4.2 metric regression suite registered in both workflows.
  • Behavioral changes worth calling out: Compared with branch-1.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa, restoring the cached subtree in the SQL graph and event log is intended. The EXPLAIN and observed-metric fallback changes are already in the supplied base. Spark 3.x can emit an additional final plan update. Observed-metric collection now reaches Spark’s scan, while the existing fallback for observed caches remains intact.
  • Suggested improvements: No further P1/P2 changes requested.

Reviewed all six files and all three commits in the full diff from d98fd2494c9a6cc8552efe6aa0a73d5714486b23 to c190f6e3ac0252987d3d77dc7bdf15313ce7147c. The PR remains non-draft. Read AGENTS.md, the snapshot and live discussions. Routed skill: .ai/skills/review-comet-pr/SKILL.md. No sibling skill applies.

No introduced P1/P2 issues found within this review. The existing P2 concern is addressed by the current override.

Validation: git diff --check passed. Spark API probes passed eight cases each on 3.5.9 and 4.1.3, covering AQE, nested caches and materialization. A byte-for-byte replay of Spark 4.2’s scope-extraction method on Spark 4.1 classes reproduced the old shuffle rejection and succeeded with the current override. The replay substitutes stable scope IDs and does not validate end-to-end accumulator execution. No full Comet build or native suite ran locally because this checkout lacks built artifacts and Maven dependencies.

Exact-head CI at 2026-10-04 16:52 UTC: 8 successful, 11 running and 14 skipped checks, with no failures. Comet execution results, including the new Spark 4.2 suite, remain pending. run-all-spark-profiles is applied. Upstream Spark SQL suites and macOS checks were skipped. This review does not establish a completed CI verdict.

@comphead
comphead added this pull request to the merge queue Oct 4, 2026
Merged via the queue into apache:main with commit f8cc631 Oct 4, 2026
57 checks passed
andygrove added a commit to andygrove/datafusion-comet that referenced this pull request Oct 5, 2026
Since apache#6577, CometInMemoryTableScanExec exposes Spark's own
InMemoryTableScanExec as its one subquery, so that Spark's UI draws the
cached plan below it. With the cache on by default, the DPP suite's
'filtering ratio policy fallback' caches its dimension table, and
checkPartitionPruningPredicate requires every subquery of an adaptive
plan to contain an AdaptiveSparkPlanExec, which Spark's scan does not.
Skip it there in every diff, since the query never runs it as a
subquery.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:scan Parquet scan / data reading bug Something isn't working run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

The SQL tab and event log lose the cached plan under CometInMemoryTableScan

2 participants