andygrove opened a new pull request, #6368:
URL: https://github.com/apache/datafusion-comet/pull/6368

   ## Which issue does this PR close?
   
   Closes #6292.
   
   ## Rationale for this change
   
   `CometExecIterator.serializeCometSQLConfs` passes the core count that sizes 
the native Tokio runtime to native code as `spark.executor.cores`. 
`numDriverOrExecutorCores` took it from the thread count of a `local[N]` 
master, or for any other master from `spark.executor.cores`, falling back to 1.
   
   A standalone (`spark://`) executor without `spark.executor.cores` takes 
every core its worker offers and runs that many tasks at once, but its 
`SparkConf` has no `spark.executor.cores`. `CoarseGrainedExecutorBackend.run` 
builds the executor's conf from the driver's `spark.*` properties and only sets 
`spark.executor.cores` for a non-default resource profile (and Spark 3.4 never 
sets it). So every such executor got a single Tokio worker, and so did 
`local-cluster` executors. Native plans that read no input from the JVM run 
entirely on Tokio workers, so all of the executor's concurrent tasks shared one 
thread (13.0 s instead of 5.5 s in the issue's measurement). On `main` this can 
also deadlock under memory pressure, when that one worker parks in Spark's 
`ExecutionMemoryPool.acquireMemory`.
   
   YARN and Kubernetes were not affected: their default for 
`spark.executor.cores` is 1, which is also the executor's number of task slots.
   
   ## What changes are included in this PR?
   
   - A new pure function, `CometExecIterator.numDriverOrExecutorCores(master, 
executorCores, availableProcessors)`, resolves the core count from the first of 
these that applies:
     - local mode: the master's thread count (unchanged)
     - `spark.executor.cores`, when set (unchanged)
     - `local-cluster[N, C, M]`: `C`, the cores each worker offers
     - `spark://`: `Runtime.getRuntime.availableProcessors()`. A standalone 
worker offers all of its machine's processors by default 
(`WorkerArguments.inferDefaultCores`), and without `spark.executor.cores` the 
master launches one executor per worker with all of its free cores 
(`Master.scheduleExecutorsOnWorkers`).
     - `yarn` and `k8s://`: 1, their default
     - any other master (for example `mesos://` on Spark 3.x, or an external 
cluster manager): unknown
   - For an unknown master, the `SparkConf` wrapper falls back to 1 as before 
and logs a warning once per JVM that points at `spark.executor.cores` and 
`COMET_WORKER_THREADS`. The warning is limited to this case rather than every 
count of 1 on a cluster, because one worker is right for a one-core executor 
(an explicit `spark.executor.cores=1`, or the YARN and Kubernetes default), and 
warning on each of those executors would be noise.
   - `spark.task.cpus` needs no handling: every task takes that many cores, so 
the core count is at least the number of task slots.
   - `COMET_WORKER_THREADS` still overrides the count; the native side is 
unchanged.
   - The "Configuring Tokio Runtime" section of the tuning guide describes the 
new resolution and what happens when an executor runs more tasks at once than 
it has worker threads.
   
   A standalone worker started with a core count other than its machine's 
(`SPARK_WORKER_CORES`) still gets the machine's processor count. Fewer offered 
cores only leave spare threads. For more, `spark.executor.cores` or 
`COMET_WORKER_THREADS` sets the count, as the tuning guide says.
   
   Open PR #6300 rewrites the worker-thread bullet in 
`docs/source/contributor-guide/development.md` to say that Comet uses one 
worker when `spark.executor.cores` is not set outside local mode. This PR 
leaves that file alone, so the bullet will need a follow-up once both have 
landed. Related to #6261, which reduces the deadlock risk independently of this 
change.
   
   ## How are these changes tested?
   
   - New test `the native runtime gets a worker thread for every core the 
executor runs tasks on` in `CometExecIteratorLifecycleSuite` covers `local`, 
`local[N]`, `local[*]` and `local[N, F]` (including `local[N]` with 
`spark.executor.cores` set, which local mode ignores), `local-cluster` and 
`spark://` (single and HA master URLs) with and without `spark.executor.cores`, 
YARN and Kubernetes with and without it, and an unknown master. The test calls 
a function this PR introduces, so there is no red run against the old code, 
which resolved `spark://` and `local-cluster` without `spark.executor.cores` to 
1.
   - On the default Spark 4.1 profile, ran the new test and the two `SQLConf 
serde` tests in `CometExecSuite`, which go through the `SparkConf` wrapper. All 
passed.
   - Checked the Spark behavior the fix relies on against the 3.4.3, 4.1.3 and 
4.2.0 sources: executors receive every `spark.*` property of the driver, 
including `spark.master` (`CoarseGrainedSchedulerBackend.sparkProperties`); 
`spark.executor.cores` is set on the executor only for a non-default resource 
profile, and never on 3.4; a standalone executor gets its cores through 
`--cores`; and the `local-cluster` master URL format.
   


-- 
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