comphead commented on code in PR #6368:
URL: https://github.com/apache/datafusion-comet/pull/6368#discussion_r4146856184
##########
spark/src/main/scala/org/apache/comet/CometExecIterator.scala:
##########
@@ -680,26 +680,69 @@ object CometExecIterator extends Logging {
}
}
+ private val unknownCoresWarned = new AtomicBoolean(false)
+
+ /**
+ * The number of cores this JVM runs tasks on, which sizes the native
runtime's worker threads.
+ * Falls back to 1, with a warning logged once, if the master does not tell.
+ */
private def numDriverOrExecutorCores(conf: SparkConf): Int = {
- def convertToInt(threads: String): Int = {
- if (threads == "*") Runtime.getRuntime.availableProcessors() else
threads.toInt
+ val master = conf.get("spark.master")
+ numDriverOrExecutorCores(
+ master,
+ conf.getOption("spark.executor.cores"),
+ Runtime.getRuntime.availableProcessors()).getOrElse {
+ if (unknownCoresWarned.compareAndSet(false, true)) {
+ logWarning(
+ s"Comet cannot tell how many cores this executor has from
spark.master=$master when " +
+ "spark.executor.cores is not set, so its native runtime starts one
worker thread. " +
+ "Native plans that read no input from the JVM run on the worker
threads, so " +
+ "concurrent tasks share this one. Set spark.executor.cores, or the
" +
+ "COMET_WORKER_THREADS environment variable, to the executor's core
count. " +
+ s"${CometConf.TUNING_GUIDE}.")
+ }
+ 1
}
+ }
- // If running in local mode, get number of threads from the spark.master
setting.
- // See
https://spark.apache.org/docs/latest/submitting-applications.html#master-urls
- // for supported formats
+ // Master URL formats, as SparkContext parses them. See
+ //
https://spark.apache.org/docs/latest/submitting-applications.html#master-urls
+ // `local[*]` means using all available cores and `local[2]` means using 2
cores.
+ private val LOCAL_N_REGEX = """local\[([0-9]+|\*)\]""".r
+ // `local[num-worker-threads, max-failures]`
+ private val LOCAL_N_FAILURES_REGEX =
"""local\[([0-9]+|\*)\s*,\s*([0-9]+)\]""".r
+ // `local-cluster[num-workers, cores-per-worker, memory-per-worker-mib]`
+ private val LOCAL_CLUSTER_REGEX =
+ """local-cluster\[\s*([0-9]+)\s*,\s*([0-9]+)\s*,\s*([0-9]+)\s*]""".r
- // `local[*]` means using all available cores and `local[2]` means using 2
cores.
- val LOCAL_N_REGEX = """local\[([0-9]+|\*)\]""".r
- // Also handle format `local[num-worker-threads, max-failures]
- val LOCAL_N_FAILURES_REGEX = """local\[([0-9]+|\*)\s*,\s*([0-9]+)\]""".r
+ /**
+ * The number of cores the driver (in local mode) or an executor runs tasks
on, given the master
+ * URL, `spark.executor.cores` and the number of processors available to the
JVM. Every task
+ * takes `spark.task.cpus` of these cores, so this is at least the number of
tasks that run at
+ * once. None if it cannot be determined.
+ */
+ def numDriverOrExecutorCores(
+ master: String,
+ executorCores: Option[String],
+ availableProcessors: Int): Option[Int] = {
+ def localThreads(threads: String): Int =
+ if (threads == "*") availableProcessors else threads.toInt
- val master = conf.get("spark.master")
master match {
- case "local" => 1
- case LOCAL_N_REGEX(threads) => convertToInt(threads)
- case LOCAL_N_FAILURES_REGEX(threads, _) => convertToInt(threads)
- case _ => conf.get("spark.executor.cores", "1").toInt
+ // Local mode runs tasks on the master's threads and ignores
spark.executor.cores.
+ case "local" => Some(1)
+ case LOCAL_N_REGEX(threads) => Some(localThreads(threads))
+ case LOCAL_N_FAILURES_REGEX(threads, _) => Some(localThreads(threads))
+ case _ if executorCores.isDefined => executorCores.map(_.toInt)
+ // Without spark.executor.cores, a standalone executor takes every core
its worker offers,
+ // but Spark does not set spark.executor.cores on the executor. A
local-cluster worker
+ // offers the master's cores per worker, and a standalone worker offers
all of the
+ // machine's processors unless it is started with a different number.
+ case LOCAL_CLUSTER_REGEX(_, coresPerWorker, _) =>
Some(coresPerWorker.toInt)
+ case _ if master.startsWith("spark://") => Some(availableProcessors)
+ // The default of spark.executor.cores on YARN and Kubernetes.
+ case _ if master == "yarn" || master.startsWith("k8s://") => Some(1)
+ case _ => None
Review Comment:
For an unknown master this still starts a single worker, which is the
situation in #6292. For `mesos://` on Spark 3.x, an executor without
`spark.executor.cores` takes the cores of the offer
(`MesosCoarseGrainedSchedulerBackend.executorCores` in v3.5.8), which is
similar to standalone. Would `availableProcessors` be a safer fallback than 1
here? I'd expect a few idle Tokio threads to cost less than starving concurrent
tasks, but I haven't measured that or looked at other external cluster managers.
##########
docs/source/user-guide/latest/tuning.md:
##########
@@ -38,11 +38,22 @@ guide is split into the following pages:
## Configuring Tokio Runtime
-Comet uses a global tokio runtime per executor process. By default it starts
one worker thread per executor core
-(`spark.executor.cores`, or the thread count of `local[N]` and `local[*]`
masters) and allows up to 512 blocking
-threads, which is tokio's default. If `spark.executor.cores` is not set
outside local mode, Comet starts a single
-worker thread. These values can be overridden using the environment variables
`COMET_WORKER_THREADS` and
-`COMET_MAX_BLOCKING_THREADS`.
+Comet uses a global tokio runtime per executor process. By default it starts
one worker thread per core that the
+executor runs tasks on, and allows up to 512 blocking threads, which is
tokio's default. Comet takes the number of
+cores from the first of these that applies:
+
+- the thread count of a `local`, `local[N]` or `local[*]` master, in local mode
+- `spark.executor.cores`, when it is set
+- the cores per worker `C` of a `local-cluster[N, C, M]` master
+- the number of processors available to the executor's JVM on a standalone
(`spark://`) cluster, since a standalone
+ executor without `spark.executor.cores` takes every core its worker offers,
and a worker offers all of its
+ machine's cores unless it is started with a different number
(`SPARK_WORKER_CORES`)
+- one on YARN and Kubernetes, which is their default for `spark.executor.cores`
+
+On any other cluster manager, Comet starts a single worker thread when
`spark.executor.cores` is not set, and logs a
Review Comment:
Now that #6300 has merged, the worker threads bullet in
`docs/source/contributor-guide/development.md` says Comet starts one worker
thread when `spark.executor.cores` is not set outside local mode. With this
change that is no longer true for `spark://` and `local-cluster`. Would it be
worth updating that bullet in this PR so it does not land stale?
##########
spark/src/main/scala/org/apache/comet/CometExecIterator.scala:
##########
@@ -680,26 +680,69 @@ object CometExecIterator extends Logging {
}
}
+ private val unknownCoresWarned = new AtomicBoolean(false)
+
+ /**
+ * The number of cores this JVM runs tasks on, which sizes the native
runtime's worker threads.
+ * Falls back to 1, with a warning logged once, if the master does not tell.
+ */
private def numDriverOrExecutorCores(conf: SparkConf): Int = {
- def convertToInt(threads: String): Int = {
- if (threads == "*") Runtime.getRuntime.availableProcessors() else
threads.toInt
+ val master = conf.get("spark.master")
+ numDriverOrExecutorCores(
+ master,
+ conf.getOption("spark.executor.cores"),
+ Runtime.getRuntime.availableProcessors()).getOrElse {
+ if (unknownCoresWarned.compareAndSet(false, true)) {
+ logWarning(
+ s"Comet cannot tell how many cores this executor has from
spark.master=$master when " +
+ "spark.executor.cores is not set, so its native runtime starts one
worker thread. " +
Review Comment:
Nit: the native side prefers `COMET_WORKER_THREADS` when it is set, so this
warning, and the single worker it mentions, would be misleading for executors
that already followed its advice. Would it make sense to skip the warning when
that variable is set?
##########
spark/src/main/scala/org/apache/comet/CometExecIterator.scala:
##########
@@ -680,26 +680,69 @@ object CometExecIterator extends Logging {
}
}
+ private val unknownCoresWarned = new AtomicBoolean(false)
+
+ /**
+ * The number of cores this JVM runs tasks on, which sizes the native
runtime's worker threads.
+ * Falls back to 1, with a warning logged once, if the master does not tell.
+ */
private def numDriverOrExecutorCores(conf: SparkConf): Int = {
- def convertToInt(threads: String): Int = {
- if (threads == "*") Runtime.getRuntime.availableProcessors() else
threads.toInt
+ val master = conf.get("spark.master")
+ numDriverOrExecutorCores(
+ master,
+ conf.getOption("spark.executor.cores"),
+ Runtime.getRuntime.availableProcessors()).getOrElse {
+ if (unknownCoresWarned.compareAndSet(false, true)) {
+ logWarning(
+ s"Comet cannot tell how many cores this executor has from
spark.master=$master when " +
+ "spark.executor.cores is not set, so its native runtime starts one
worker thread. " +
+ "Native plans that read no input from the JVM run on the worker
threads, so " +
+ "concurrent tasks share this one. Set spark.executor.cores, or the
" +
+ "COMET_WORKER_THREADS environment variable, to the executor's core
count. " +
+ s"${CometConf.TUNING_GUIDE}.")
+ }
+ 1
}
+ }
- // If running in local mode, get number of threads from the spark.master
setting.
- // See
https://spark.apache.org/docs/latest/submitting-applications.html#master-urls
- // for supported formats
+ // Master URL formats, as SparkContext parses them. See
+ //
https://spark.apache.org/docs/latest/submitting-applications.html#master-urls
+ // `local[*]` means using all available cores and `local[2]` means using 2
cores.
+ private val LOCAL_N_REGEX = """local\[([0-9]+|\*)\]""".r
+ // `local[num-worker-threads, max-failures]`
+ private val LOCAL_N_FAILURES_REGEX =
"""local\[([0-9]+|\*)\s*,\s*([0-9]+)\]""".r
+ // `local-cluster[num-workers, cores-per-worker, memory-per-worker-mib]`
+ private val LOCAL_CLUSTER_REGEX =
+ """local-cluster\[\s*([0-9]+)\s*,\s*([0-9]+)\s*,\s*([0-9]+)\s*]""".r
- // `local[*]` means using all available cores and `local[2]` means using 2
cores.
- val LOCAL_N_REGEX = """local\[([0-9]+|\*)\]""".r
- // Also handle format `local[num-worker-threads, max-failures]
- val LOCAL_N_FAILURES_REGEX = """local\[([0-9]+|\*)\s*,\s*([0-9]+)\]""".r
+ /**
+ * The number of cores the driver (in local mode) or an executor runs tasks
on, given the master
+ * URL, `spark.executor.cores` and the number of processors available to the
JVM. Every task
+ * takes `spark.task.cpus` of these cores, so this is at least the number of
tasks that run at
+ * once. None if it cannot be determined.
+ */
+ def numDriverOrExecutorCores(
+ master: String,
+ executorCores: Option[String],
+ availableProcessors: Int): Option[Int] = {
+ def localThreads(threads: String): Int =
+ if (threads == "*") availableProcessors else threads.toInt
- val master = conf.get("spark.master")
master match {
- case "local" => 1
- case LOCAL_N_REGEX(threads) => convertToInt(threads)
- case LOCAL_N_FAILURES_REGEX(threads, _) => convertToInt(threads)
- case _ => conf.get("spark.executor.cores", "1").toInt
+ // Local mode runs tasks on the master's threads and ignores
spark.executor.cores.
+ case "local" => Some(1)
+ case LOCAL_N_REGEX(threads) => Some(localThreads(threads))
+ case LOCAL_N_FAILURES_REGEX(threads, _) => Some(localThreads(threads))
+ case _ if executorCores.isDefined => executorCores.map(_.toInt)
+ // Without spark.executor.cores, a standalone executor takes every core
its worker offers,
+ // but Spark does not set spark.executor.cores on the executor. A
local-cluster worker
+ // offers the master's cores per worker, and a standalone worker offers
all of the
+ // machine's processors unless it is started with a different number.
+ case LOCAL_CLUSTER_REGEX(_, coresPerWorker, _) =>
Some(coresPerWorker.toInt)
+ case _ if master.startsWith("spark://") => Some(availableProcessors)
+ // The default of spark.executor.cores on YARN and Kubernetes.
+ case _ if master == "yarn" || master.startsWith("k8s://") => Some(1)
Review Comment:
Nit: once this is rebased, `isContainerSizedFromOverhead(master)` (added in
#6375) tests the same `yarn` or `k8s://` condition. Would it make sense to
share one check, perhaps under a neutral name, so the two cannot drift apart?
The intent differs (container memory versus default cores), so keeping both is
also reasonable.
--
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]