andygrove opened a new issue, #6292: URL: https://github.com/apache/datafusion-comet/issues/6292
### Describe the bug `CometExecIterator.numDriverOrExecutorCores` (`CometExecIterator.scala:641-662`) sizes the Tokio runtime from the `local[N]` master, or for any other master from `spark.executor.cores`, falling back to 1. A standalone (`spark://`) executor with `spark.executor.cores` unset takes every core its worker offers and runs that many tasks at once. Its SparkConf still has no `spark.executor.cores`, because Spark only sets that on the executor for a non-default resource profile (4.1.3 `CoarseGrainedExecutorBackend.scala:484-505`). So the whole executor gets one Tokio worker. `local-cluster` masters behave the same way. YARN and Kubernetes are fine, because their default of one core is also the number of task slots. A plan with no JVM input runs entirely on Tokio workers (`jni_api.rs:1152-1189`). With one worker, every such task on the executor shares a single thread. Measured on `main` at `634e37d08` with a native scan feeding `sortWithinPartitions` over 4.8M rows (`local[4]`, 2g off-heap, `COMET_WORKER_THREADS` standing in for the missing core count): | | 1 worker | 4 workers | |---|---|---| | `main` | 13.0 s | 5.5 s | | with #6261 | 7.7-8.6 s | 5.5 s | #6261 hands a worker's core to another thread for the duration of each memory acquire, which hides part of the cost. I'd expect the gap to grow with the executor's core count. On `main` this can also deadlock. With 96m or 128m of off-heap memory, the same query hung in 4 runs out of 4. The only worker was parked in Spark's `ExecutionMemoryPool.acquireMemory`, waiting for 1/2N of the pool, and the task holding that memory could only release it by running on that same worker. #6261 fixes this: 11 runs out of 11 passed, including one where the wait happened and cleared. The tuning guide documents the one-worker fallback (`tuning.md:42-45`), but not its cost. ### Steps to reproduce Run a stage whose native plan has no JVM input, for example a native Parquet scan feeding a sort, on a standalone cluster without `spark.executor.cores`. The executor log shows `Comet tokio runtime: using spark.executor.cores=1 worker threads`. Locally, `COMET_WORKER_THREADS=1` with `local[4]` reproduces the same numbers. ### Expected behavior The runtime gets at least one worker per task slot. ### Additional context For `spark://` masters without `spark.executor.cores`, `Runtime.getRuntime.availableProcessors()` matches what a standalone worker gives the executor by default. A warning when the resolved count is 1 on a non-local master would also help. `COMET_WORKER_THREADS` works around it today. -- 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]
