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]