Skip to content

fix: keep operators above a cached relation native after the AQE re-plan - #6208

Merged
andygrove merged 2 commits into
apache:mainfrom
andygrove:fix-issue-6202
Sep 25, 2026
Merged

andygrove merged 2 commits into
apache:mainfrom
andygrove:fix-issue-6202

Conversation

@andygrove

@andygrove andygrove commented Sep 24, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6202.

Rationale for this change

On Spark 3.5+, AQE wraps each cache scan in a TableCacheQueryStageExec and, once the stage materializes, plans the operators above it again. CometExecRule had no case for the stage, so those operators stayed on Spark. For SELECT k, count(*) FROM t GROUP BY k over a cached t, the initial plan is fully native, but the final plan runs both aggregates and the exchange on Spark over a CometColumnarToRow. Results are correct, so the cost is performance, and it becomes more visible once #5634 turns the native cache scan on by default.

What changes are included in this PR?

Two cases in CometExecRule, one per cache format:

  • A query stage whose plan is a CometInMemoryTableScanExec is a native input through CometExchangeSink, as a Comet shuffle or broadcast stage already is. The scan produces Arrow batches, and foreachUntilCometInput already accepts any QueryStageExec as an input.
  • Spark's cache scan does not come back from the re-plan as a bare stage. When the CometScanWrapper around a CometSparkToColumnarExec is removed, the CometSparkToColumnarExec inherits the scan's logical link, so AQE reuses it over the stage as the relation's physical plan. It is not a CometNativeExec, so the operators the re-plan puts above it could not convert. An existing CometSparkToColumnarExec is now wrapped in a CometScanWrapper again, the way isCometScan nodes already are on every pass.

Neither case applies with default settings, since the native cache scan and spark.comet.sparkToColumnar.enabled are both off.

How are these changes tested?

A new test in CometInMemoryCacheSuite runs the three queries from the issue (an aggregate, a broadcast join and a filtered aggregate) with the native cache on and off, each first over a cold cache and then over the warm one, and checks that the executed final plan has no Spark operators. The check is a new CometTestBase.checkCometOperatorsInFinalPlan, which runs checkCometOperators over the final plan and the plan inside each query stage. checkSparkAnswerAndOperator could not catch this, because it inspects an adaptive plan that has not run, so it only ever sees the initial plan. The cold run uses checkToRDD = false, since checkAnswer otherwise runs the query once as an RDD first and warms the cache. That suite registers Comet's extensions twice (through spark.plugins and withExtensions), so there the re-plan brings back two stacked conversions. A second test in CometExecSuite covers a single registration with Spark's default cache format. Both fail without the change.

Locally, the cache suites pass on Spark 3.4 (where the new tests are skipped), 3.5, 4.0, 4.1 and 4.2, and CometExecSuite, CometAggregateSuite, CometJoinSuite, the shuffle suites and CometExecRuleSuite pass on 4.1.

On Spark 3.5+, AQE wraps each cache scan in a TableCacheQueryStageExec and,
once the stage materializes, plans the operators above it again. CometExecRule
had no case for the stage, so those operators stayed on Spark in the final
plan even though the initial plan was fully native.

A stage over a Comet in-memory table scan is now a native input, as a Comet
shuffle stage is. Spark's cache scan comes back from the re-plan inside the
CometSparkToColumnarExec converted over it, which inherits the scan's logical
link, so an existing CometSparkToColumnarExec is wrapped in a CometScanWrapper
again to let the operators above it convert.
@github-actions github-actions Bot added the bug Something isn't working label Sep 24, 2026
Move the final AQE plan check into CometTestBase as
checkCometOperatorsInFinalPlan, which runs checkCometOperators over the
final plan and the plan inside each query stage, and share
isTableCacheStage between the two suites.

checkAnswer runs the query once as an RDD before collecting it, which
warmed the cache before the checked run. Run each query cold and then
warm on one cache, with checkToRDD = false.

@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: AQE replanning could leave aggregates, joins and exchanges above cached relations on Spark despite an initially native plan.
  • Design approach: Recognize stages containing CometInMemoryTableScanExec through CometExchangeSink, and rewrap existing CometSparkToColumnarExec nodes with CometScanWrapper.
  • Correctness / compatibility analysis: Checked Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The changes fit Spark’s stage execution and logical-link reuse. Traced schema checks, Arrow input dispatch, partitioning and wrapper removal without finding an introduced correctness issue.
  • Key design decisions: Reusing existing planning wrappers keeps the change small. The wrappers are removed before execution and introduce no additional runtime conversion or batch-copy mechanism. Matching QueryStageExec avoids a dependency on the table-cache stage class absent from Spark 3.4.
  • Implementation sketch: Two planner cases, a shared final-plan checker, and regression coverage for cold/warm caches, native-cache enablement and single extension registration.
  • Behavioral changes worth calling out: Eligible cached queries retain native parent operators after AQE replanning. Configuration defaults remain unchanged. Performance improvements were not independently benchmarked.
  • Suggested improvements: No introduced P1/P2 issues found within this review.

Reviewed both commits and all four changed files against base 02e84047a423a08debe3c6b6daab19fd9233427c, at full head a9b257b388e183196a77a61b667f5d15af96018c. The PR remained non-draft. Snapshot and live discussion checks contained no existing reviews, comments or threads.

Routed skills: review-comet-pr and review-comet-ffi-pr for the Arrow input boundary.

Exact-head CI: Required Checks passed. The Spark 4.1/JDK 17 execution job passed both new regression tests and reported 1,021 successful tests with zero failures. Native tests, scans, expressions, shuffle and TPC-H/TPC-DS checks also passed.

Validation limits: No local build or runtime tests were run because this checkout lacks the built native library and JVM dependencies. Spark SQL, Iceberg and macOS jobs were skipped, and other Spark runtime profiles were not exercised by this CI run. The author’s reported cross-version local results were not independently reproduced. This review does not establish the broader Spark SQL verdict required before queueing a planner change.

@andygrove
andygrove added this pull request to the merge queue Sep 25, 2026
@andygrove

Copy link
Copy Markdown
Member Author

Thanks @sunchao

Merged via the queue into apache:main with commit a096c9e Sep 25, 2026
38 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Operators above a cached relation fall back to Spark after AQE materializes the table-cache stage

2 participants