Skip to content

Comet starts a single Tokio worker on standalone executors when spark.executor.cores is unset #6292

Description

@andygrove

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

bugSomething isn't workingperformancepriority:mediumFunctional bugs, performance regressions, broken features

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions