Repository navigation
feat: add spark.comet.explain.planOnly.enabled - #5394
Conversation
Builds the Comet plan Comet would have executed and logs it to the driver log, then discards it and lets Spark run the query unchanged. Both CometScanRule and CometExecRule short-circuit when the config is set; CometExecRule uses thread-local bypass flags on both rules to force normal behavior while it constructs the preview for the report. The WARN is deduped per SQL execution ID so AQE emits one report per query. Closes apache#5335.
- Collapse two ThreadLocals into one shared `planOnlyPreviewInProgress` (dropped `CometScanRule.forceApply`/`withForceApply`). - Move the plan-only branch to the top of the exec-enabled block so `normalizePlan`/`RewriteJoin`/`tagUnsafePartialAggregates` are not run and thrown away. - Replace hand-rolled `LinkedHashSet` LRU with a synchronized `LinkedHashMap` + `removeEldestEntry`, matching `IcebergPlanDataInjector.commonCache`. Drops the `clearPlanOnlyReported` test helper. - Collapse V1/V2 × AQE test variants into one loop; fold `assertNoComet` into `runPlanOnlyAndAssertReverted` and drop redundant `collect`s.
`_apply` on both rules now takes/uses an explicit `forPreview` flag rather than reading a thread-local. `reportPlanOnlyCoverage` calls the private `_apply` methods directly, skipping the `Rule.apply` short-circuits, so the recursive preview reads the parameter instead of ambient state. Widens `CometScanRule._apply` to `private[rules]` so `CometExecRule` can invoke it. Drops `planOnlyPreviewInProgress` and `withPreview`.
sunchao
left a comment
There was a problem hiding this comment.
Summary
Reviewed eb514de1ff76785b81a4940d6a2adb3ea149ab07 against a28ac348f5a68500d96fe872b49e453d96e4bab6 across five independent scopes. Reusing the existing conversion rules is a sensible way to keep eligibility checks aligned with normal Comet planning. I found two P2 issues that affect the accuracy and completeness of the report. Both are described inline.
Prior state and problem
The existing explain output describes a plan that has already been converted for Comet execution. Evaluating an unfamiliar workload therefore requires enabling Comet on the execution path. This change adds an opt-in way to inspect potential acceleration while retaining Spark execution.
Design approach
The new configuration defaults to false. The public scan rule leaves the incoming plan alone, while the execution rule builds a disposable preview by invoking scan conversion and execution conversion directly. The explicit forPreview argument prevents that preview from taking the plan-only shortcut again.
Correctness / compatibility analysis
I did not find a verified native-execution leak in the ordinary read path. The main concerns are reporting correctness: nested subquery planning can claim the execution-ID report slot before the outer query, and the preview stops before Scala-side transition insertion and transition-driven reversion.
A local Spark 3.5.2 planning-order probe using the same deduplication logic reproduced the first issue with AQE both off and on. In both cases the scalar aggregate was reported first and the outer filter was suppressed. The second issue is supported by the exact-head rule ordering and the existing RevertNativeForTransitionHeavyStagesSuite regression case. I did not run the full Comet JVM/native suite. git diff --check passed. At the final CI check, 12 checks had succeeded, 5 were running, 1 was queued, 7 were skipped, and none had failed.
Key design decisions
Keeping the feature disabled by default limits ordinary behavior changes. Returning the original plan avoids reconstructing a Spark plan from converted operators. A bounded execution-ID cache also limits retained reporting state, but the report owner must distinguish the root query from recursively prepared subqueries.
Implementation sketch
The implementation adds the configuration and user guide, exposes scan conversion within the rules package, adds the preview/report branch to CometExecRule, and tests V1/V2 scans with AQE on and off plus a scalar subquery and a config-off sanity check. Those tests check plan isolation, but they do not inspect the warning emitted by an actual query action.
Behavioral changes worth calling out
Plan-only mode requires Comet execution to be enabled, emits a driver warning, and deliberately does not call DataFusion's native planner. That documented native-planning limitation is reasonable. It should remain distinct from omitting a known Scala-side reversion or reporting only an inner subquery.
Suggested improvements
Please make the root query own the once-per-execution report, and calculate coverage after the applicable columnar-transition and Comet post-columnar rules. Add action-based log assertions for scalar subqueries under both AQE modes and compare the preview with normal execution when transition reversion is enabled.
| if (!forPreview && CometConf.COMET_EXPLAIN_PLAN_ONLY_ENABLED.get()) { | ||
| val executionId = Option( | ||
| session.sparkContext.getLocalProperty(SQLExecution.EXECUTION_ID_KEY)) | ||
| if (CometExecRule.markPlanOnlyReported(executionId)) { |
There was a problem hiding this comment.
[P2] Keep the report slot for the outer query
Could we avoid consuming the execution-ID entry while Spark is preparing a nested subquery? For a fresh action on SELECT id FROM range(10) WHERE id > (SELECT max(id) FROM range(3)), Spark prepares the scalar subquery before the outer plan under the same execution ID. The subquery records the ID here, so the later outer-query call is suppressed. I reproduced this planning order with the same deduplication logic on Spark 3.5.2 with AQE both disabled and enabled. The only report describes max(id), not the outer scan/filter or the workload being evaluated. Please distinguish root-query reporting from subquery/stage preparation and add an action-based test that captures the warning and checks that the outer plan appears.
| * both rules run their normal transforms instead of short-circuiting. | ||
| */ | ||
| private def reportPlanOnlyCoverage(plan: SparkPlan): Unit = { | ||
| val preview = _apply(CometScanRule(session)._apply(plan), forPreview = true) |
There was a problem hiding this comment.
[P2] Include post-columnar reversion in the coverage preview
Could the preview run transition insertion and Comet's post-columnar rules before computing coverage? Normal planning subsequently runs RevertNativeForTransitionHeavyStages and EliminateRedundantTransitions, but this path stops before them. With spark.comet.exec.transitionRevert.enabled=true, spark.comet.exec.transitionRevert.maxTransitions=0, and Comet project execution disabled, the existing regression case shows that the real executed plan has zero CometExec nodes. This preview still counts the operators that the configured Scala-side rule removes, and its transition count is calculated before Spark inserts those transitions. That is separate from the documented uncertainty about DataFusion planning failures. Please compare the report with the real post-columnar plan for this configuration.
Plan-only reporting keyed its once-per-query slot on the SQL execution ID alone. Spark prepares scalar and DPP subqueries as their own top-level plans before the outer plan reaches the conversion rules, so a nested subquery took the slot and the outer query - the plan being evaluated - was never reported. Reports are now keyed on the execution ID and the plan, so the outer query always gets one, while AQE's per-stage and per-re-optimization applications of the rule are recognised as re-plans and stay quiet. The preview also stopped at operator conversion, before Spark inserts columnar transitions and Comet runs its post-columnar rules. It therefore counted operators that RevertNativeForTransitionHeavyStages removes and counted transitions before they existed. The preview now inserts transitions and applies both post-columnar rules, using a new applyToAllStages entry point so that every shuffle boundary is visited rather than only the topmost stage. Tests capture the logged report and assert that the outer query appears under both AQE modes, and that the reported coverage matches the coverage of the plan Comet really executes when stage reversion fires.
|
Both fixed in 2f94707. On the first one — I couldn't find a way to identify the root plan at rule-application time. Subqueries are prepared before the outer plan with nothing to distinguish them, and matching the incoming plan against the root QueryExecution's output breaks down for commands and write plans, which would then get no report at all. So instead of one report owned by the root, each independently planned plan gets one: the outer query plus one per separately prepared subquery. Dedupe is now keyed on execution ID and plan, and AQE's per-stage and post-re-optimization applications are skipped as re-plans of something already reported, so the stage spam the execution-ID cache was there to prevent is still gone. Does that seem like a reasonable trade to you, or would you rather see the root identified even if some plan shapes drop out of reporting? The second one was as you described. The preview now inserts transitions and runs both post-columnar rules. RevertNativeForTransitionHeavyStages needed a new entry point for this: its What the preview still can't match is AQE re-planning: it describes the pre-adaptive plan and applies the post-columnar rules to it in one pass rather than per stage. That's called out in the user guide alongside the native-planning caveat. |
| df.collect() | ||
| val executed = CometCoverageStats.forPlan(df.queryExecution.executedPlan) | ||
|
|
||
| val reports = withSQLConf(CometConf.COMET_EXPLAIN_PLAN_ONLY_ENABLED.key -> "true") { |
There was a problem hiding this comment.
[P1] Preserve the captured report across Spark 3.x withSQLConf
Could the withSQLConf block be moved inside capturePlanOnlyReports, or could these assertions run inside the configuration block? Spark 3.4 and 3.5 define withSQLConf(...)(f: => Unit): Unit, unlike the generic Spark 4.x helper, so reports is inferred as Unit here. The exact-head Spark 3.4 and Spark 3.5 checks both fail test compilation on the subsequent size, mkString, and head calls. This prevents the supported Spark 3.x test builds from compiling.
| plan: SparkPlan, | ||
| queryStagePrep: Boolean): Boolean = { | ||
| executionId match { | ||
| case None => |
There was a problem hiding this comment.
[P2] Deduplicate adaptive reports when the execution ID is absent
Could the no-ID path retain plan-scoped reporting state instead of returning true for every invocation? The public df.rdd.count() path can build and execute AQE stages without installing spark.sql.execution.id. A Spark 3.5.2 probe using this exact decision logic and both rule registrations produced five report decisions for SELECT id % 2 AS k, count(*) AS n FROM range(20) GROUP BY id % 2: initial preparation, the adaptive wrapper, the exchange stage, adaptive re-optimization, and the final stage. A fresh collect() produced one. Because this branch bypasses both the stage check and deduplication, plan-only mode rebuilds previews and emits overlapping coverage summaries for those ordinary RDD-backed workloads, contrary to the documented suppression of stage/re-optimization reports. Please cover df.rdd.count() and planning via executedPlan before an action in the reporting tests.
| * holds the whole plan at once, whereas under AQE Spark hands that rule one stage at a time. | ||
| */ | ||
| private def reportPlanOnlyCoverage(plan: SparkPlan): Unit = { | ||
| val converted = _apply(CometScanRule(session)._apply(plan), forPreview = true) |
There was a problem hiding this comment.
[P2] Include converted scalar subqueries in the outer coverage report
Even with AQE disabled, an acceleratable scalar subquery is counted as Spark in the outer report. Its independently converted preview has already been discarded, and these two conversion passes walk ordinary plan children, leaving the original ScalarSubquery.plan in the outer preview. ExtendedExplainInfo then traverses those expression-owned plans and includes their operators in the percentage. For the new scalar-subquery test query, the separate subquery warning can therefore report acceleration that is missing from the outer query's coverage. A constructed-plan probe using the exact-head formatter/serializer and real Comet project nodes reports 1/5 with the untouched subquery versus 2/5 after replacing only its plan. Please carry the converted subquery previews into the outer preview, or exclude separately reported subqueries from that report's counts, and compare its coverage with normal Comet planning.
|
Seems like there are CI failures :) |
…on ID, preview subqueries Move the coverage assertions inside their withSQLConf block. Spark 3.4 and 3.5 declare withSQLConf as returning Unit, so binding its result inferred Unit and broke test compilation on every supported Spark 3.x build. Replace the execution-ID keyed report slots with a mark on the plan itself. df.rdd.count() and reading executedPlan without an action can plan, and in the first case execute AQE stages, with no execution ID installed, and the old no-ID path reported unconditionally. Catalyst copies tags onto replacement nodes, so the mark survives the rewrites between one rule application and the next and identifies an AQE stage or final-plan pass as a re-plan of something already reported. The execution-scoped hash is kept for the one case a mark cannot cover: a subquery referenced twice is prepared twice, as separate but identical plans. Preview the plan behind each subquery expression before computing coverage. Extended explain counts the plans owned by a node's expressions, and normal planning has already converted those by the time the outer plan arrives, so the outer report counted accelerated subquery operators as Spark. Two things were needed to make the numbers line up: a prepared subquery arrives with codegen wrappers and transitions already inserted, which blocks conversion, so those are stripped first; and ReuseExchangeAndSubquery is replayed at the top level, since it is the last preparation step and otherwise each copy of a subquery is counted separately. With AQE off the outer report now equals CometCoverageStats for the plan Comet really executes. Under AQE the outer report still describes the pre-adaptive plan, in which subqueries have not been planned yet and so count as Spark. Called out in the user guide.
|
All three fixed in 1683a65. The Spark 3.x compile break was exactly as you described — the assertions now run inside the On the no-execution-ID path: rather than scoping state to an ID that may not exist, the report is now marked on the plan itself. Catalyst copies tags onto replacement nodes, so the mark survives the rewrites Spark makes between one application of the rule and the next, which is what lets a later application recognise a stage or final-plan pass as a re-plan of something already reported. That turned out to matter beyond the no-ID case: I found a plan with no exchanges gets a final-plan pass under AQE whose hash no longer matches the prep-time plan, because the subquery was substituted in between, so the old dedupe let it through as a second report. I kept the execution-scoped hash alongside the mark for the one case a mark can't cover — a subquery referenced from two places is prepared twice as separate but identical plans. Both The subquery one was worse than the double-count you found, and getting the outer report to agree with normal planning took two more things. A subquery reaches the outer plan fully prepared, so its aggregates are already inside One thing I could not make match, and I'd like your read on it. Under AQE the outer report is produced from the prep-rule application, and at that point The remaining red check is |
sunchao
left a comment
There was a problem hiding this comment.
Re-reviewed 1683a65 with five agents. The 41-test CometExecRuleSuite and 96 normal/preview plan comparisons passed on Spark 4.1.3. Targeted probes against locally compiled JVM and native code reproduced the two reporting issues below; the third comment identifies a test coverage gap.
| // Reuse bookkeeping: the plan to preview is one level further down. | ||
| case reused: ReusedSubqueryExec => reused.copy(child = previewSubquery(reused.child)) | ||
| case other => | ||
| val preview = buildPreview(stripPreparation(other.child), topLevel = false) |
There was a problem hiding this comment.
[P2] Apply stage reversion to the prepared DPP child
Could the DPP broadcast child be previewed through its own preparation boundary before restoring the exchange wrapper? With AQE disabled, Spark's PlanDynamicPruningFilters prepares the child before constructing BroadcastExchangeExec. This path instead passes that exchange into buildPreview, so applyToAllStages treats it as a boundary and leaves the reconverted child unreverted. I reproduced this on the PR head with a partitioned Parquet fact table joined to a filtered Parquet dimension, spark.comet.exec.transitionRevert.enabled=true, spark.comet.exec.transitionRevert.maxTransitions=0, and spark.comet.exec.project.enabled=false: the outer report says 4/12 accelerated (33%), while CometCoverageStats for the normally executed plan says 2/12 (16%). Both executions returned the same ten correct rows. The separate DPP report correctly shows 0/3, but its child is counted as accelerated again in the outer report. Please preserve the child's normal preparation scope and add a non-AQE DPP coverage comparison.
| plan.exists(p => | ||
| p.isInstanceOf[QueryStageExec] || p.isInstanceOf[AdaptiveSparkPlanExec] || | ||
| p.getTagValue(PLAN_ONLY_REPORTED).isDefined) || | ||
| (aqeEnabled && !queryStagePrep && plan.isInstanceOf[Exchange]) |
There was a problem hiding this comment.
[P2] Keep report state when AQE replaces an empty plan
Could report ownership survive replacement of the entire physical tree? With plan-only mode and AQE enabled, spark.sql("SELECT id % 2 AS k, count(*) AS n FROM range(20) WHERE id < 0 GROUP BY id % 2").collect() emits an 83% report for the initial aggregate and then a second 0% report for EmptyRelation. After the empty shuffle materializes, AQE creates a fresh plan containing neither QueryStageExec nor the reporting tag. Its structural hash also differs, so both guards admit another report under the same execution ID. I reproduced the two actual Comet warnings on Spark 4.1.3; AQE disabled emits one. Please retain query/subquery reporting ownership across this rewrite and cover an empty adaptive query.
| } { | ||
| val label = s"${if (useV1) "V1" else "V2"} scan, AQE=$aqe" | ||
| test(s"plan-only mode: $label") { | ||
| withParquetTable((0 until 100).map(i => (i, i % 5)), "tbl") { |
There was a problem hiding this comment.
[P3] Select V2 before creating the Parquet fixture
Could USE_V1_SOURCE_LIST be configured outside withParquetTable? The fixture reads Parquet and registers tbl before runPlanOnlyAndAssertReverted changes that setting. Changing it afterward does not replace the existing logical relation, so both tests labeled "V2 scan" still exercise FileSourceScanExec. A Spark 4.1.3 probe confirmed this with AQE on and off; creating the fixture after the config change produces BatchScanExec. Actual V2 plan-only execution passed separate checks, but these committed tests do not cover that path.
…coverage Preview a DPP subquery inside its BroadcastExchangeExec rather than around it. PlanDynamicPruningFilters prepares the build plan and wraps it in the exchange afterwards, so the plan that went through the post-columnar rules is the exchange's child. Handing the exchange to the preview left that child a stage whose top boundary was the exchange, so RevertNativeForTransitionHeavyStages never fired and the outer report claimed 4/12 accelerated where the executed plan has 2/12. Suppress reports for the rest of an execution once AQE has cut a query stage. When a shuffle materializes empty, AQE re-plans from the logical plan and can collapse the tree to an empty relation; that plan shares no nodes with the one already reported, so the mark does not reach it, and it holds no query stages either, so the structural checks do not either. Its hash differs too, so both guards admitted a second 0% report under the same execution ID. Subqueries are all compiled before the first stage is cut, so this costs no report except for a subquery AQE only plans once stages are under way, which the user guide now notes. Set USE_V1_SOURCE_LIST before creating the Parquet fixture in the plan-only scan tests. withParquetTable resolves the relation through spark.read and registers the result as a temp view, so changing the source list afterwards left both the V1 and the V2 variant planning a FileSourceScanExec. The scan node is now asserted so this cannot regress unnoticed. New tests cover an empty adaptive query under AQE on and off, and compare the outer report for a DPP subquery against CometCoverageStats for the plan Comet really executes with stage reversion forced on.
|
All three fixed in c75a87c. The DPP one was as you described. The empty adaptive plan I reproduced too, and it turned out the hash was not the only guard that missed: the mark could not reach the new tree either, since You were right about the V2 tests, and the cause is that Separately, while building the DPP fixture I hit an unrelated pre-existing bug and filed #5486: with AQE on, DPP in play, and
|
sunchao
left a comment
There was a problem hiding this comment.
Re-reviewed c75a87c with five agents. The 60 existing tests across five suites and 132 non-AQE coverage comparisons passed on Spark 4.1.3. Targeted probes reproduced the two reporting issues below.
| executionId.foreach(planOnlyStagedExecutions.add) | ||
| false | ||
| } else if (isReapplication(plan)) { | ||
| false | ||
| } else if (executionId.exists(planOnlyStagedExecutions.contains)) { |
There was a problem hiding this comment.
[P2] Suppress empty AQE reports without an execution ID
The new guard only retains state when an execution ID exists. With AQE and spark.comet.explain.planOnly.enabled=true, an empty grouped query through the JVM bridge used by PySpark's df.rdd still emits the initial 83% report followed by a misleading 0% EmptyRelation report:
SELECT id % 2 AS k, count(*) AS n
FROM range(20)
WHERE id < 0
GROUP BY id % 2I reproduced this on the current head with Spark 4.1.3 by invoking javaToPython() and counting the resulting RDD; AQE disabled emits one report. The action has no SQL execution ID, and the replacement tree has no reporting tag. Could suppression survive this AQE replacement without requiring an execution ID, with a regression covering this RDD path?
| // Node tags are not part of a plan's structural hash, so the mark set above does not | ||
| // perturb the key. Without an execution ID there is nothing to scope the state to, and the | ||
| // checks above have already ruled out the repeat applications AQE makes, so report. | ||
| executionId.forall(id => planOnlyReportedPlans.add(s"$id:${plan.hashCode()}")) |
There was a problem hiding this comment.
[P3] Canonicalize nested subqueries before deduplicating reports
Projecting the same nested scalar subquery twice and calling .collect() with AQE disabled produces four reports under one execution ID: the inner aggregate, the outer aggregate twice, and the main query. For a Parquet table tbl with integer columns _1 and _2, the reproducer is:
SELECT _1,
(SELECT max(_2) FROM tbl
WHERE _1 > (SELECT min(_2) FROM tbl)) AS a,
(SELECT max(_2) FROM tbl
WHERE _1 > (SELECT min(_2) FROM tbl)) AS b
FROM tblOn Spark 4.1.3 at the current head, the two outer-aggregate reports are byte-for-byte identical, and both normal and plan-only physical plans reuse that subquery through ReusedSubqueryExec. The raw plan hash misses this equivalence, causing redundant previews and warnings. Could the key use a canonical fingerprint, with a nested-subquery reuse regression? This is a reporting issue; results and coverage values still match.
…reused subqueries Mark the logical plan a reported physical plan was built from, alongside the existing mark on the physical tree. When a stage materializes empty AQE re-plans from the logical plan and can collapse the whole tree to an empty relation, leaving a physical tree that shares no node with the one reported and holds no query stage, so only the logical mark reaches it. The mark needs no SQL execution ID, which the previous stage-tracking guard did, so the duplicate 0% report is now suppressed on the RDD path PySpark's `df.rdd` takes as well. That also removes the reason the stage-tracking guard existed, and with it the suppression of every report that followed the first stage cut, so a subquery AQE only plans once stages are under way - a DPP subquery - gets its report back. Key the per-execution dedupe on the canonicalized plan. A subquery referenced twice is prepared once per reference as plans differing only in expression IDs, which Spark then collapses to one `ReusedSubqueryExec`; the raw structural hash saw two plans and reported the same work twice.
sunchao
left a comment
There was a problem hiding this comment.
[P2] Re-reviewed 53db476c. The original aggregate/no-ID case is addressed, but the existing empty-AQE reporting issue still has a source-verified residual. When Spark removes an outer local sort before Comet's preparation rule, Comet marks the surviving logical Join, while AQE retains the original, unmarked logical Sort. Empty propagation replaces the Join and then the Sort; the replacement root inherits the Sort's unmarked tags and can pass the reporting guards again without a SQL execution ID.
The additional trigger is a two-partition Range merge join on a.id % 7 = b.id % 7, filtered by a.id < 0, with SORT BY a.id % 7 and execution via queryExecution.toRdd.count(). The added aggregate/no-ID and nested-subquery regressions do not cover this removed-root shape. This path was traced through current source and Spark 4.1.3, not executed in a runtime reproduction.
A mark on the plan already reported cannot reach the plan AQE builds when a stage materializes empty, and marking the logical plan only reached it while the physical root's logical link happened to be the logical root. When a preparation rule drops the root operator first - `RemoveRedundantSorts` removing a sort above a sort merge join - the mark lands on the join below it, empty propagation replaces the join and then the sort, and the replacement root inherits the unmarked sort's tags. Recognize that plan by what it is instead. A re-optimized plan holds a `QueryStageExec` for every stage that materialized; the only way to eliminate all of them is empty propagation, which leaves a plan Catalyst records as producing zero rows, so the logical link answers the question directly. The check is confined to the query-stage-prep rule, so a genuinely empty query is still reported. This subsumes the logical mark, which is dropped.
|
I am moving this to draft. This seems overly complex. I will explore alternative approaches. |
The interpolator broke both scalafix lint jobs. The guard keeps a preview failure from failing the query, which is the one thing plan-only mode promises not to do.
`IcebergWriteStrategy` is registered with `injectPlannerStrategy`, not as a `Rule[SparkPlan]`, so it runs during physical planning ahead of both conversion rules and their plan-only short-circuit never reaches it. An Iceberg V2 write was still emitted as Comet's two-operator shape and executed natively, which is exactly what plan-only mode promises not to do. Neither `IcebergCommitExec` nor `IcebergWriteExec` extends `CometPlan`, so the plan-only suite's "no Comet operators" assertion could not see it. Decline the strategy in plan-only mode. The cost is that the report no longer describes the write as accelerated; that is documented alongside the existing native-planner caveat. Everything else registered in `CometSparkSessionExtensions.apply` is already inert once no Comet nodes exist: `CometPlanAdaptiveDynamicPruningFilters` only matches Comet scans, `CometReuseSubquery` creates no Comet nodes, and the two post-columnar rules only act on Comet operators.
The two were registered as separate rules, adjacently, in both the columnar and the query-stage-prep paths. Nothing ever ran between them, and neither is useful alone: CometExecRule seeds its native chain only from the nodes CometScanRule produces, so operator conversion over unconverted Spark scans converts nothing. Compose them so the ordering is an invariant of the code rather than of the registration order, and so callers that need the whole conversion have one entry point. Both rules keep their classes, files and tests.
|
The reason that this is so complex is because we have two separate rules - CometScanRule and CometExecRule. If these get combined into a single rule then we can just create the comet plan and discard it and no longer need to roll things back. Working on it .... |
Stacks this branch on the rule composition (apache#6082) and simplifies plan-only mode to use it: - CometScanRule goes back to main's version. It no longer knows about plan-only mode and its `_apply` is private again. - The `forPreview` parameter is gone from CometExecRule, along with the plan-only branch in `_apply` and the `queryStagePrep` constructor parameter. CometExecRule is now purely operator conversion. - The plan-only state, the reporting decision and the preview machinery move to CometRule, which is the one place that holds the whole conversion. `buildPreview` calls `convert(...)` instead of reaching into both rules' private `_apply`s. Plan-only mode is scoped to `spark.comet.exec.enabled` in `planOnlyApplies`, matching where the check used to sit inside the exec-enabled branch, so a plan the conversion rules would have left alone is not diverted into a report and columnar shuffle keeps being applied with exec disabled.
|
Rebased onto #6082, which composes The reason for splitting it out is that most of the awkwardness in this PR came from the two rules being registered separately. With one rule there is a single place that holds the whole conversion, so plan-only mode gets a single short-circuit and a single call to build the plan it reports on:
One behavioural note. The plan-only check used to sit inside Also fixed here: Local runs: |
# Conflicts: # spark/src/main/scala/org/apache/comet/CometSparkSessionExtensions.scala # spark/src/main/scala/org/apache/comet/rules/CometRule.scala # spark/src/test/scala/org/apache/comet/rules/CometExecRuleSuite.scala
|
rebased on main |
|
@sunchao on the removed-root-sort residual from your last review: this is handled in 0896e36. Rather than chasing the mark through the rewrite, the query-stage-prep rule now recognizes a re-plan that AQE collapsed to nothing directly: its logical link has |
Share the post-columnar rule list between CometColumnar and the plan-only preview through CometRule.postColumnarRules, replacing the hand-synced copy. RevertNativeForTransitionHeavyStages takes a wholePlan flag in place of the separate applyToAllStages entry point. Remove the unused imports left in CometExecRule, which failed the scalafix lint, and restore the CometRule docstring and the "scan conversion must run before operator conversion" test from apache#6082 that the main merge dropped. Condense the comments, the user guide section and the config doc, inline single-use helpers, and factor the plan-only tests onto shared helpers. Check the Iceberg split-write conf before the plan-only conf so the common path does one lookup.
|
Reviewed
The overall design is sound: composing conversion in Validation: Current CI has 23 successful checks and 14 skipped. Its execution group passed 999 tests, including all 62 Local probes used Spark 4.1.3/JDK 17 with freshly compiled current Scala changes over cached dependencies and native library. Disk constraints prevented a clean full build. These are targeted integration reproductions; no whole-query performance claim is made. |
Render plan-only reports in the verbose format regardless of spark.comet.explain.format, so the coverage summary is never dropped. Tell AQE re-planning apart from initial preparation by whether the rule is running inside AdaptiveSparkPlanExec.reOptimize, instead of guessing from maxRows. An initially empty query is now reported under AQE, and a re-plan that keeps operators above an empty join is no longer reported twice. Deduplicate repeated subqueries within the query being prepared on the current thread rather than within a SQL execution ID, so queries planned without one (toRdd, df.rdd, executedPlan) no longer repeat reports.
|
@sunchao all three addressed in 3123088.
These checks read Spark method names off the stack. I checked that |
|
@sunchao could you take another look? I'd like to get this one into the 1.1.0 release. Thanks! |
|
I'm going to go ahead and merge this. It is an optional feature, disabled by default, and I want to create the 1.1 branch tomorrow. Feel free to create follow on issues and assign them to me if there is more feedback! |
Since apache#5394 the post-columnar rules also run in plan-only mode, where Comet only reports the plan it would execute and Spark executes its own. CometCacheColumnarRule still inserted a ColumnarToRowExec there. Return early when plan-only mode applies, using the predicate CometRule uses, except in the plan-only preview, which should still show the transition. In the runtime-settings test, add a plan-only case and use QueryTest.checkAnswer with checkToRDD = false so the cold iteration really reads a cold cache.
* perf: reduce Spark cache row conversion overhead * chore: remove cache row reader benchmark results * perf: feed cached Arrow columns into Spark codegen * docs: illustrate Comet cache columnar rewrite * fix: honor Comet disable switches for fused cache reads * fix: align cache fusion with Spark codegen settings * ci: retry after runner and dependency download failures * perf: split the generated cache row reader for wide projections The generated CachedBatchRowIterator passed every column read to GenerateUnsafeProjection through ctx.currentVars, which disables its method splitting, so next() held one write per column. Past about 100 nullable columns it exceeded HotSpot's 8000-byte HugeMethodLimit and was never JIT-compiled, and near 1500 columns it failed to compile at all. Read each column through a codegen-only leaf expression instead, so GenerateUnsafeProjection splits the field writes into bounded methods as it does for any projection, and bind a batch's columns in a loop. The largest generated method is now 371 bytes at 200 nullable bigint columns and 1905 bytes at 1500 (was 16462 bytes and a compile failure). Also check the ByteCodeStats that CodeGenerator.compile returns, as WholeStageCodegenExec does: when a method exceeds min(spark.sql.codegen.hugeMethodLimit, 8000), read the batches with UnsafeProjection over batch.getRow in an indexed loop without copy(). The interpreted path uses the same loop with InterpretedUnsafeProjection. * fix: do not fuse cache reads in plan-only mode Since apache#5394 the post-columnar rules also run in plan-only mode, where Comet only reports the plan it would execute and Spark executes its own. CometCacheColumnarRule still inserted a ColumnarToRowExec there. Return early when plan-only mode applies, using the predicate CometRule uses, except in the plan-only preview, which should still show the transition. In the runtime-settings test, add a plan-only case and use QueryTest.checkAnswer with checkToRDD = false so the cold iteration really reads a cold cache. * test: check fused cache reads with Comet on and native execution off The row-consumer test ran with Comet disabled, so no case took the fused path and the test could not notice. Run it with Comet on and native execution and shuffle off, with AQE on and off, add a filter that reads every column through codegen, and assert per query whether the fused transition is present. * test: measure the fused cache reader in CometInMemoryCacheBenchmark Fold the fused-consumer case into CometInMemoryCacheBenchmark's Spark-operator section instead of a separate harness. Its Comet-off arm now measures the row reader, so name it that, add an arm with Comet on and native execution off (on-heap enabled, so Comet loads and the read fuses), and check each arm's reader in its plan. Add wide relations of 100, 200 and 1500 columns read through the row readers. * docs: describe how Spark operators read Comet's cache format The Limitations section still said Spark-operator reads were slower for reasons not yet established, with a table measured through the row path this PR replaced. Explain the three ways a Spark operator reads a relation cached in Comet's format (under native execution, through the fused ColumnarToRowExec and when it applies, and through the row reader otherwise) and replace the table with one that measures both readers. * fix: avoid numeric widening in cache benchmark * ci: retry preflight after Maven Central download failure --------- Co-authored-by: Andy Grove <agrove@apache.org>
Which issue does this PR close?
Closes #5335. Alternate approach to #5345, based on review discussion there.
Heads up: I used an LLM to help draft this. The design is mine, but the code and prose have been shaped with LLM assistance, so review with that in mind.
Rationale for this change
Users evaluating Comet on a workload need a way to estimate how much of it Comet would accelerate without actually changing execution. Turning Comet on and comparing runs carries real risk. This adds a mode that builds the Comet plan Comet would have executed, logs it, and lets Spark run the query unchanged.
What changes are included in this PR?
spark.comet.explain.planOnly.enabled, default off.CometScanRuleandCometExecRuleshort-circuit at the top of theirapplyand return the plan untouched.CometExecRulebuilds a preview by calling both rules' private_applymethods directly (with an explicitforPreview = trueparameter that skips the short-circuit on the recursive call), logs the resulting Comet plan at WARN level, then returns the original Spark plan.normalizePlan/RewriteJoin/tagUnsafePartialAggregatesare not run and thrown away.understanding-comet-plans.md.The estimate is Scala-side only. The native plan is never handed to DataFusion, so anything that would have failed inside DataFusion still counts as accelerated. That is called out in the config docstring and the user guide.
How are these changes tested?
New tests in
CometExecRuleSuitecover:Each test asserts the executed plan has zero
CometPlanoperators.