apache / apache/datafusion-comet

Do not accelerate a query stage that contains a columnar-to-row fallback (enable whole-stage revert by default)

Open
#5,207 0 comments 0 reactions 1 assignee Claimed by @andygrove View on GitHub
enhancement performance user experience
Dominant language
Scala
Stars
1.3k
Forks
373
Avg merge
2d 6h
Merged PRs (30d)
190

Description

## What is the problem the feature request solves?

When someone tries Comet for the first time on a query where Comet accelerates the scan and then immediately falls back to Spark row-based execution, the result is often *slower* than plain Spark: the columnar-to-row (C2R) conversion has to be paid, and the surrounding Spark operators lose whole-stage codegen because a Comet operator sits underneath them. This is a bad first impression, and the user has no obvious signal that anything is wrong — the plan looks like Comet "worked".

Comet already has a whole-stage fallback for this, but it is disabled by default and, as currently written, cannot be enabled at the strictest setting.

### What exists today

`RevertNativeForTransitionHeavyStages` (`spark/src/main/scala/org/apache/comet/rules/RevertNativeForTransitionHeavyStages.scala`), registered in `postColumnarTransitions` ahead of `EliminateRedundantTransitions`, counts `ColumnarToRowTransition` nodes within a stage (stopping at `QueryStageExec` / shuffle / broadcast boundaries) and, above a threshold, rewrites every `CometExec` back to `cometExec.originalPlan`, strips the Comet transitions, re-inserts Spark's, and wraps in `RowToColumnarExec` when the parent exchange needs columnar input. AQE applies it per query stage; non-AQE applies it via `transformUp` over shuffles plus the result stage.

- `spark.comet.exec.transitionRevert.enabled` — default `false`
- `spark.comet.exec.transitionRevert.maxTransitions` — default `2`

Adjacent but much narrower: `spark.comet.shuffle.revertRedundantColumnar.enabled` (default `true`) only unwinds a Comet columnar shuffle sandwiched between two non-Comet hash aggregates.

## Describe the potential solution

### 1. Exempt output-boundary C2R transitions from the count

`countTransitions` currently counts the trailing C2R at the plan root, so a *fully* accelerated query such as `CometNativeScan -> CometProject -> ColumnarToRowExec(root)` counts as 1 transition. Setting `maxTransitions=0` today would therefore revert every query, including ones with no mid-stage fallback at all.

A C2R at the output boundary is not a fallback — it is the unavoidable handoff to the consumer. `countTransitions` should skip any C2R reachable from the stage root through only output-only nodes:

- the stage root itself
- `DataWritingCommandExec`
- `WriteFilesExec`
- `V2CommandExec` (covers append / overwrite / CTAS)
- `ExecutedCommandExec`

Every other C2R still counts.

This exemption is a prerequisite for the default change below, not an optional refinement.

### 2. Change the defaults

- `spark.comet.exec.transitionRevert.enabled` -> `true`
- `spark.comet.exec.transitionRevert.maxTransitions` -> `0`

That is: Comet does not attempt to accelerate a stage at all if the stage contains any non-exempt C2R fallback.

### 3. Make the explain output actionable

When the rule fires it records `withFallbackReason(reverted, s"Stage reverted: $transitionCount C2R transitions exceed threshold $maxTransitions")`, which does surface in `spark.comet.explain.format=verbose` / `=fallback` and (when enabled) in `spark.comet.explain.fallback.log.enabled` warnings. Per-operator fallback reasons survive the revert too, since untranslated Spark nodes keep their own tags. But there are gaps that matter a lot more once this is on by default:

- **The message never names the config**, so a user cannot discover the off switch. Compare `COMET_EXPLAIN_FALLBACK_ENABLED`, which embeds `set =false` directly in its own log line (`CometExecRule.scala:610`). The revert reason should name `spark.comet.exec.transitionRevert.enabled` (and the `maxTransitions` knob) inline.
- **The one warning that fires on the discovery path cannot see the revert.** `spark.comet.explain.fallback.enabled` logs from inside `CometExecRule`, i.e. `preColumnarTransitions` — strictly before the revert rule runs. It needs to run (or be duplicated) after `postColumnarTransitions` so the reverted stage is reported.
- `spark.sql.extendedExplainProviders` only exists on Spark 4.0+ and is not set by default, so on Spark 3.4/3.5 the annotated plan is only reachable by calling `ExtendedExplainInfo` directly. Meanwhile the coverage summary will now read "Comet accelerated 0 out of N eligible operators" with no visible cause. Worth considering whether a one-time `logWarning` should be unconditional when a stage is reverted.

### 4. Docs

Update the `CometConf` docs for both configs, regenerate `docs/source/user-guide/latest/configs.md` via `GenerateDocs`, and rewrite the "Reducing Row/Columnar Conversion Overhead" section of `docs/source/user-guide/latest/tuning.md` (currently documents the old defaults).

## Trade-offs and expected fallout

Reverting the whole stage also gives up the Comet Parquet scan, which is often the single largest win. `SELECT my_udf(a), b FROM t WHERE b > 5` currently keeps `CometNativeScan + CometFilter`; under the new default it becomes all-Spark. The premise of this change is that the C2R cost plus the loss of whole-stage codegen usually outweighs the scan win — that should be confirmed with benchmarks before merging, and `maxTransitions=1` is the hedge if it does not hold.

Test fallout to plan for:

- TPC-DS / TPC-H plan-stability golden files across all profile directories (`approved-plans-v1_4`, `-spark3_5`, `-spark4_0`, `-spark4_1`, `approved-plans-v2_7`, ...) will shift and need regenerating.
- Any suite asserting *partial* acceleration needs `spark.comet.exec.transitionRevert.enabled=false` pinned explicitly.

## Additional context

Targeting after the 1.0 release, since this is a user-visible default behavior change with broad plan-stability impact.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.