Skip to content

Experimental local execution mode that runs a whole query as one DataFusion graph #6499

Description

@viirya

What is the problem the feature request solves?

When Spark runs with a local master, the driver and all executors share one JVM, but Comet still executes a query the distributed way: each Spark task builds and runs its own native plan, stages exchange data through Spark shuffle files, and each result partition is encoded by its task and decoded on the driver. In a single process these boundaries are pure overhead. Partial aggregates, join inputs and sort runs are written to and read back from shuffle files that DataFusion could exchange in memory, and the native plans of one stage cannot share operator state.

DataFusion can already execute a whole physical plan, including hash repartitioning, on its own Tokio runtime. For queries that Comet fully supports, the single-process case could hand the whole query to DataFusion and keep Spark only for parsing, analysis, optimization and result delivery.

Describe the potential solution

Add an experimental, opt-in local execution mode (an internal spark.comet.exec.local.enabled option, off by default):

  • Admit a query only as a whole, before ordinary Comet conversion: Spark 4.1, an in-process local master, AQE disabled, batch queries without subqueries, and a fully supported plan. Any other query keeps the existing Comet/Spark path, and there is no fallback once native output has started.
  • Plan the whole query as one DataFusion graph in the driver, reusing Comet's native Parquet scan and Spark-compatible expressions. Replace Spark exchanges with DataFusion repartitioning inside that graph, so the query has no Spark shuffle and runs as a single Spark result task. Each execution owns a fresh graph and a query-scoped memory pool with spilling.
  • Start with Parquet scan/filter/project, grouped and global COUNT/MIN/MAX, shuffled hash joins, and terminal global sort, Top-K and limit, and leave everything else to the existing path.
  • Keep Spark's cancellation, job groups and SQL metrics by still running the result task as a Spark job, while delivering collect/take results to the driver within the same JVM instead of encoding them in the task.

Additional context

Related to #1204, which avoids repeating native plan construction across tasks within an executor. This mode instead targets the single-process case, where the whole query can be one native graph.

Activity

  1. viirya commented on Oct 4, 2026

    @viirya
    MemberAuthor

    Thanks @andygrove for the careful review on #6500. Replying here since the main questions are about direction.

    Reusing the existing operator translation. Agreed. local.proto and the local planner re-describe aggregate, join, sort and limit and lower them a second time, so every newly admitted operator would need a second implementation without the Spark-compatibility handling that PhysicalPlanner::create_plan already has. I'd rather build the graph from the Operator trees that CometExecRule produces:

    • Admit a plan only when it is entirely native: native blocks connected by CometShuffleExchangeExec, with no JVM input.
    • Lower each block with PhysicalPlanner::create_plan, plan all file partitions of a native scan in the one graph, and replace each shuffle exchange with a RepartitionExec, or for range partitioning with a per-partition sort and a sort-preserving merge.
    • Admission then follows whatever Comet runs natively, and local.proto goes away.

    One thing I'd like to confirm: within one graph both sides of a join go through the same RepartitionExec, so I don't think the local exchange needs Comet's Spark-compatible Murmur3 partitioning, unless some native operator depends on it.

    AQE. Admission from the CometRule(session, queryStagePrep = true) instance looks like the right way to drop the session-wide requirement. Since InsertAdaptiveSparkPlan leaves exchange-free plans alone, those can be admitted with AQE enabled from the columnar rule, and plans with exchanges from the query stage prep rule. I'll prototype the latter to check that replacing the initial plan with a single leaf leaves AQE nothing to re-plan.

    Splitting. Happy to split it along the lines you suggested:

    1. The native/local lifecycle, JNI bridge, result delivery and range.
    2. Admission of exchange-free native plans (scan/filter/project) from CometExecRule operator trees, without requiring AQE to be disabled.
    3. Exchanges through RepartitionExec for aggregates and joins, admitted from the query stage prep rule, plus Top-K and limit.
    4. Global sort after the DataFusion 56 upgrade ([EPIC] DataFusion 56 upgrade: regressions found by tracking DataFusion main #6410); see the review thread on feat: Experimental local execution mode that runs a whole query as one DataFusion graph #6500.

    The other review points (spill directories and size limit, native partition count, unordered limit determinism, cached plans, the Spark version gate and memory_management.md) would be addressed in the PR that introduces the affected code. I'll convert #6500 to a draft and keep it as a reference until the first split PR is up.

  2. andygrove commented on Oct 7, 2026

    @andygrove
    Member

    Thanks, the plan looks good to me.

    On the Murmur3 question: two RepartitionExecs with DataFusion's hash do co-partition with each other, so the hash only matters when one side of a join gets its partitioning from a leaf instead of an exchange. Bucketed tables are the common case. CometNativeScanExec reports the bucket HashPartitioning, and EnsureRequirements keeps the side that needs no exchange and shuffles only the other side to match it. If that exchange becomes a RepartitionExec with DataFusion's hash, partition i of the shuffled side no longer lines up with bucket i and the join returns wrong rows. Could admission decline plans where a join child's partitioning comes from a scan (bucketed HashPartitioning, or KeyGroupedPartitioning from a storage-partitioned join), or could the local exchange keep Comet's Spark-compatible hash? A join between a bucketed and a non-bucketed table on the bucket column would make a good test.

    Lowering each block once with PhysicalPlanner::create_plan has a related catch. Every DataFusion partition shares one instance of each planned expression, and that instance carries the planner's partition index. So spark_partition_id(), monotonically_increasing_id(), the rand and randn seeds, Sample's seed and the partition index JvmScalarUdfExpr hands the dispatcher would all see partition 0. Could the deterministic gate and the nativeOnly protobuf walk in LocalParquetPlanner carry over to the operator-tree admission?

    For the first PR, the range needs an overflow guard. The local range admits ranges whose arithmetic overflows in Spark. On this head, spark.range(Long.MinValue, Long.MaxValue, 1L << 62, 1) returns four rows with local execution and none without, because Spark's generated code overflows. spark.range(Long.MinValue, Long.MaxValue, 1, 1).take(3) returns three rows against none. Main now has a native RangeScan (#6553) whose mayOverflow check keeps these ranges on Spark. Could the first PR lower RangeScan through create_plan instead of adding a second range operator, and compare the range tests against Spark's answer rather than Scala ranges?

    The session has a similar gap. query_context starts from SessionContext::new_with_config_rt, which registers DataFusion's default functions and none of Comet's overrides. That's harmless while nativeOnly keeps ScalarFunc out. Once admission follows CometExecRule trees, concat would resolve to DataFusion's version, which skips nulls, instead of SparkConcat. Could the local session come from the per-task setup (prepare_datafusion_session_context, or create_execution_context once #6695 lands)? That would also bring the spill directories and their size cap, the spark.comet.datafusion.* pass-through and the skip-partial-aggregation guard.

    Last, which budget should bound a local query? Each one gets its own 256 MiB FairSpillPool that Spark never sees, so concurrent local queries add up with nothing bounding the total, and the spark.memory.offHeap.size users already set for Comet doesn't apply. The result task is a real Spark task. Could the pool come from its CometTaskMemoryManager, as in the per-task path, so spark.comet.exec.local.memoryLimit isn't needed?

  3. andygrove commented on Oct 7, 2026

    @andygrove
    Member

    @viirya once this is working, it should be easy to have it run against Ballista too! cc @avantgardnerio

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    enhancementNew feature or request

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions