viirya opened a new pull request, #6500: URL: https://github.com/apache/datafusion-comet/pull/6500
## Which issue does this PR close? Closes #6499. ## Rationale for this change With a local master, Comet still runs a query as Spark stages: per-task native plans, Spark shuffle files between stages, and per-task result encoding. Within one process these boundaries are overhead. This PR adds an experimental, opt-in mode that runs a whole admitted query as one DataFusion graph in the driver, so DataFusion exchanges data in memory and Spark keeps only planning, cancellation and result delivery. ## What changes are included in this PR? Local execution is enabled by the internal `spark.comet.exec.local.enabled` option (default false). See `docs/source/contributor-guide/local-execution.md` for the full design. - **Admission (`spark-local`)**: `CometLocalRule` runs before ordinary Comet conversion and admits a query only as a whole: Spark 4.1, `SparkContext.isLocal`, AQE disabled, batch queries without subqueries. Admitted shapes are Parquet (DataSource V1, `file:` paths) scan/filter/project, grouped and global `COUNT`/`MIN`/`MAX`, a single shuffled hash join, terminal global sort, Top-K and limit, and range projections. Expressions and types are limited to an explicit allowlist. Everything else keeps the existing Comet/Spark path; there is no fallback after native output starts. - **Native planning (`native/core/src/local/planner.rs`)**: builds one graph per execution, reusing core's `PhysicalPlanner` for Parquet scans and expressions. Aggregation runs as partial aggregate, DataFusion hash repartition and final aggregate. Joins repartition both sides into a partitioned `HashJoinExec`. Global sort sorts each partition and merges with `SortPreservingMergeExec`. `local.proto` describes the aggregation, join and terminal operations. - **Execution lifecycle (`native/local`, `native/core/src/local.rs`)**: query-owned graph, session config, runtime environment and `FairSpillPool` (`spark.comet.exec.local.memoryLimit`, spill controlled by `spark.comet.exec.local.spill.enabled`); a bounded one-batch handoff to the JVM; a JNI registry of numeric query IDs with cancellation that does not need the reader lock. - **Result delivery**: `CometLocalResultExec` runs the single result task as a Spark job (keeping cancellation, job groups and SQL metrics), but `collect`/`take` copy rows to the driver within the same JVM instead of encoding and compressing them in one task. It enforces `spark.driver.maxResultSize` while copying (on uncompressed `UnsafeRow` bytes, so more conservatively than Spark). - **Sort memory settings**: the planner sizes `sort_spill_reservation_bytes` and `sort_in_place_threshold_bytes` by the number of sorters sharing the budget. The in-place threshold works around a DataFusion 55.1 `ExternalSorter` bug, fixed in DataFusion 56.0.0, where a spilling sort can fail to grow a new unspillable merge reservation instead of spilling. The code comment and docs mark it for removal after upgrading. - **Benchmark**: `dev/bench-local-execution.py` and `CometLocalExecutionBenchmark` compare Spark, existing Comet and local execution, check that all results match, and include a memory-pressure check. Results are in `docs/source/contributor-guide/local-execution-benchmark.md`. - `CometRule` invokes local admission first; `CometConf` adds the local options. No existing behavior changes when the option is off. Benchmark results on five million fact rows, a 512 MiB budget and `local[4]` (medians in ms, forward / reverse mode order; all 210 executions matched): | Case | Spark | Comet | Local | |---|---:|---:|---:| | scan-filter-project | 76.2 / 83.2 | 65.5 / 70.7 | 68.1 / 65.9 | | grouped-count-min-max | 129.4 / 140.8 | 79.4 / 83.1 | 66.2 / 63.6 | | partitioned-join | 264.4 / 287.1 | 110.3 / 115.3 | 56.5 / 57.7 | | top-k | 84.3 / 81.2 | 67.1 / 56.6 | 38.0 / 39.3 | | full-sort | 869.4 / 978.9 | 710.2 / 754.6 | 366.3 / 359.1 | These are small-sample, warm-cache results on one machine. No TPC-H or TPC-DS query is admitted as a whole with the current operator surface. ## How are these changes tested? - `CometLocalExecutionSuite` (46 tests, Spark 4.1, registered in the Linux and macOS PR workflows) compares results with Spark for every admitted shape and requires a local node so that fallback cannot hide a failure. It covers fallback, repeated actions, early termination, task and job group cancellation, native errors, reservation failures, the result size limit, and checks that native query handles, imported Arrow memory and result slots return to zero after each test. - Native tests in `native/local/tests` and `native/core/src/local/planner.rs` cover exchange routing, single-use graph ownership, cancellation and teardown, spilling aggregation and sort on one Tokio worker, join reservation failures, budget isolation between concurrent queries, and a regression test for concurrent multi-column sorts spilling under a shared budget. - The benchmark's `pressure` mode verifies at 64 MiB and 128 MiB that a five-million-row sort spills and matches Spark, that it fails cleanly with spill disabled, and that a following query succeeds. - Spark 3.5 compiles with `-Pstrict-warnings`; local execution stays disabled there. The Spark SQL suite is requested through the `run-spark-4.1-tests` label. This pull request and its description were written by Isaac. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
