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]

Reply via email to