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]

Reply via email to