Skip to content

Add mode to run Comet planning but execute with Spark, so users can assess potential Comet coverage #5335

Description

@andygrove

What is the problem the feature request solves?

Users evaluating Comet have no low-risk way to find out how much of their workload Comet would actually accelerate. Enabling Comet changes execution, so any assessment carries correctness and performance risk and usually needs a separate, dedicated run.

Describe the potential solution

Add a "plan only" / dry-run mode (e.g. spark.comet.planOnly.enabled) where Comet runs its full planning pass — operator/expression conversion and fallback analysis — logs the resulting Comet plan, and then executes the original Spark plan instead.

The mechanism can reuse what spark.comet.exec.transitionRevert.enabled already does in RevertNativeForTransitionHeavyStages: after planning, replace Comet operators with their Spark equivalents. Here it would be applied unconditionally to the whole plan rather than per stage above a transition threshold.

This composes with existing reporting (extended explain fallback/coverage info and spark.comet.metrics.enabled), so a user could run their production workload unchanged and collect an estimate of Comet coverage.

Points to discuss:

  • Where to surface the Comet plan: driver log, extended explain info, and/or metrics.
  • Whether to also serialize/create the native plan, which catches native-side planning failures and gives a more accurate estimate at some planning cost.
  • Behavior under AQE, where planning happens per stage.
  • Whether this belongs as its own config or as a mode of the existing explain/telemetry features.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

enhancementNew feature or request

Type

No type

Projects

No projects

    Milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions