In Spark, how many tasks can one executor run at once, and what sets that limit?
answer
- threads in one JVM, not processes
- cores per executor, divided by something
- the divisor is usually one
- tasks per stage stays the same either way
- count waves, not machines
basics
~10 sAn executor runs spark.executor.cores divided by spark.task.cpus tasks concurrently — one task per slot, each in its own thread of the same JVM. Total application parallelism is that number times the executor count.
solid answer
~50 sEach executor is a single JVM with a fixed number of **task slots**: `spark.executor.cores` divided by `spark.task.cpus` (which defaults to 1). With `--executor-cores 4` an executor runs four tasks at the same time, each on its own thread inside that one JVM, all sharing the executor's heap. Cluster-wide concurrency is slots per executor times the number of live executors, so ten four-core executors give forty concurrent tasks. That number never changes how many tasks a stage *has* — that comes from partition count — only how many **waves** the stage takes: a 200-task stage across 40 slots runs in five waves. The two failure modes are symmetric: far more slots than tasks leaves the cluster idle, and far more tasks than slots is fine but makes each wave's slowest task set the pace. Note that a "core" here is an accounting unit Spark hands out, not an enforced CPU reservation.
code
bash · 7 linesspark-submit \
--num-executors 10 \
--executor-cores 4 \
--executor-memory 8g \
--conf spark.task.cpus=1 \
job.py
# 10 x (4 / 1) = 40 tasks run concurrentlygo deeper
Be able to do the arithmetic: executors times cores per executor equals concurrent tasks. Know that this is separate from how many tasks the stage contains.
Explain slots as threads in one executor JVM, derive wave count from task count and slot count, and say why spark.task.cpus exists and defaults to 1.
Diagnose from the Spark UI: distinguish slot starvation from skew from locality waiting, and justify an executor shape against GC behaviour and broadcast amortization.
Own the sizing policy across a shared cluster — executor shape, task-slot accounting versus real CPU use, and how the choice interacts with multi-tenant fairness and cost per job.
## Slots, not cores An executor is one JVM process. Inside it Spark maintains a thread pool and a count of **task slots**, computed as `spark.executor.cores / spark.task.cpus`. `spark.task.cpus` defaults to 1, so in practice slots equal cores. Each slot holds one running task; the scheduler will not place a fifth task on a four-slot executor until one of the four finishes. The word "core" is misleading. Spark is not reserving a physical CPU: it is handing out accounting units that the cluster manager also uses for placement. If your task's user code spawns its own threads or calls into a native BLAS library that saturates every core on the box, Spark still counts it as one slot, and the machine becomes oversubscribed. `spark.task.cpus` exists precisely for that case — set it to 2 and each task consumes two slots, halving concurrency per executor. ## The arithmetic that matters Application-wide concurrency is: `concurrent tasks = executors × (spark.executor.cores / spark.task.cpus)` And the number of **waves** a stage needs is: `waves = ceil(stage task count / concurrent tasks)` A stage with 200 tasks (a common number, since post-shuffle partitioning defaults to `spark.sql.shuffle.partitions` = 200) running on ten four-core executors takes five waves. If tasks average 30 seconds, the stage takes roughly two and a half minutes plus scheduling overhead — and it takes that long regardless of how much heap sits unused. This arithmetic is the fastest sanity check available in an interview. Given a stage's task count and a cluster's slot count, you should be able to say immediately whether the job is slot-starved, well matched, or wasting a cluster. ## The two mismatches **Fewer tasks than slots** is pure waste: a stage with 12 tasks on a cluster with 200 slots leaves 188 idle, and adding executors does nothing. The fix is upstream — more input splits, or repartitioning before the expensive stage — not more machines. **Far more tasks than slots** is normally healthy, because many small waves smooth over uneven task durations. It becomes a problem only at extremes: tens of thousands of very short tasks make per-task scheduling and serialization overhead a visible fraction of runtime, and the driver becomes the bottleneck launching and collecting them. Within any wave, the wave finishes when its slowest task finishes. That is why a single task an order of magnitude slower than its peers can stall an otherwise well-provisioned stage, and why the summary metrics table on the Spark UI stage page shows min, 25th percentile, median, 75th percentile and max task duration rather than just an average. ## Threads in one JVM, and what they share Because the slots on an executor are threads in one JVM, all concurrent tasks share the same heap, the same garbage collector and the same off-heap allocations. Increasing `spark.executor.cores` therefore increases memory pressure per executor even though it costs no extra memory setting: four tasks each spilling large hash tables into the same heap is a very different situation from four single-core executors. It also increases GC contention, which is a large part of why extremely wide executors tend to underperform in practice. They also share the executor's shuffle write and read bandwidth, and the JVM's classloader and broadcast variables — the last of which is an advantage: a broadcast table is materialized once per executor, not once per task, so wider executors amortize broadcasts better. ## Locality delays the assignment Slot availability is necessary but not sufficient for a task to launch. Spark prefers to run a task where its data already is, using locality levels in descending preference: `PROCESS_LOCAL` (data in this executor's memory), `NODE_LOCAL` (on this machine), `RACK_LOCAL`, then `ANY`. When a free slot has poor locality, the scheduler waits up to `spark.locality.wait` (default 3 seconds) for a better slot before relaxing to the next level. On a busy cluster with cached data this shows up as slots that appear free while tasks sit pending — a Spark UI symptom that looks like a scheduling bug but is deliberate. ## Where the numbers come from These are submission-time settings: `--executor-cores` / `spark.executor.cores`, `--num-executors` / `spark.executor.instances`, and `spark.task.cpus`. With dynamic allocation enabled the executor count moves at runtime, so slot count is a moving target and Spark scales executors to the pending task backlog. Whichever way they are set, the two rules stay the same: tasks come from partitions, and slots come from cores.
- When would you set spark.task.cpus above 1?When a single task legitimately uses more than one CPU — user code that spawns its own thread pool, or a native library such as a multi-threaded BLAS in an ML workload. Setting it to 2 makes each task consume two slots, so the executor runs half as many tasks and stops oversubscribing the machine. It is a blunt lever and rarely needed for plain SQL work.
- Why do very wide executors, say 32 cores each, often perform worse than several smaller ones?All 32 tasks share one heap and one garbage collector, so GC pauses lengthen and contention grows; HDFS and object-store clients also throughput-cap per JVM. The usual compromise is roughly four to five cores per executor, which keeps GC manageable while still amortizing broadcast variables across several tasks.
- A Spark UI shows free slots while tasks stay pending. What is a likely cause?Locality waiting. The scheduler prefers PROCESS_LOCAL or NODE_LOCAL placement and will hold a task for up to `spark.locality.wait` (default 3s) rather than launch it on a badly-placed free slot. It is deliberate, and lowering the wait trades data locality for faster start-up on clusters with fast networks.
saying these in an interview costs you the question
- Says executor count changes how many tasks a stage has
- Thinks each task runs in its own JVM process
- Assumes a Spark core is a reserved physical CPU
- Believes more cores per executor costs no extra memory
- Claims idle slots always mean the cluster is too small