Dynamic allocation
Add and remove executors with the workload, and keep shuffle files safe while executors come and go.
On this page
spark.dynamicAllocation.enabled, Spark asks the cluster manager for more executors when tasks queue up and releases executors that sit idle. The catch is shuffle files: they live on executors, so Spark needs an external shuffle service, shuffle tracking or decommissioning to remove executors without losing them.Static vs dynamic
With static allocation, an application holds spark.executor.instances executorsexecutor: A worker process on a cluster machine that runs tasks and holds data in memory. Learn more → from start to finish. A notebook that runs one query an hour keeps 50 executors idle in between, and a job whose last stagestage: A group of steps Spark can run without moving data between machines. A new stage starts at every shuffle. Learn more → has 10 taskstask: The work for one partition in one stage, run on one CPU core. Learn more → keeps 40 executors waiting. Dynamic allocation sizes the application to its current backlog of tasks.
How it scales
-
Backlog
Tasks have been pending for schedulerBacklogTimeout (1 s)
-
Ramp up
Request 1, then 2, 4, 8 more executors each round while the backlog persists
-
Cap
Never beyond what pending tasks need, or maxExecutors
-
Idle
An executor with no tasks for executorIdleTimeout (60 s) is released
The exponential ramp is deliberate: start cautiously in case a few executors are enough, then grow fast if they are not.
| Config | Default | Meaning |
|---|---|---|
spark.dynamicAllocation.enabled | false | Turn it on |
...minExecutors / ...maxExecutors | 0 / unlimited | Bounds. Always set a max on shared clusters. |
...initialExecutors | minExecutors | Executors at start (spark.executor.instances is used if larger) |
...schedulerBacklogTimeout | 1s | How long tasks must wait before asking for more |
...executorIdleTimeout | 60s | Idle time before an executor is released |
...cachedExecutorIdleTimeout | infinity | Idle time before an executor holding cached data is released |
The shuffle file problem
Map tasks write shuffleshuffle: Moving rows between machines so that all rows with the same key end up together. Needed by joins, groupBy and sorting, and usually the most expensive step of a job. Learn more → files to their executor's local disk, and reducers fetch them later, possibly in a later job. If an idle executor is removed, its shuffle files go with it and the stages that wrote them must re-run. Spark therefore refuses to start with dynamic allocation unless one of these is in place:
External shuffle service
spark.shuffle.service.enabled=true. A long-running service on each node (a YARN NodeManager auxiliary service, or a standalone worker service) serves shuffle files, so they outlive the executor that wrote them. The classic setup on YARN.Shuffle tracking
spark.dynamicAllocation.shuffleTracking.enabled=true (Spark 3.0+). Spark simply does not release executors that hold shuffle data still needed by an active job. No extra service; the usual choice on Kubernetes, at the cost of releasing executors later.Decommissioning
spark.decommission.enabled with spark.storage.decommission.shuffleBlocks.enabled (Spark 3.1+). Before an executor leaves, its shuffle and cached blocks are migrated to other executors or to fallback storage. Also protects against spot instance loss.When it helps, and when it hurts
| Good fit | Poor fit |
|---|---|
| Shared clusters and notebooks with idle time between queries | Short, latency-sensitive jobs: requesting and starting executors takes seconds to minutes |
| Jobs whose stages vary a lot in task count | Jobs that cache heavily: cached executors are never released by default |
| Many applications competing for one cluster | Long-running streaming queries with steady load, where a fixed size is simpler and predictable |
Common mistakes
- Enabling it without a shuffle service or tracking — Spark refuses to start the application.
- No maxExecutors on a shared cluster — One big job takes every executor the cluster manager will give.
- Caching and expecting executors to scale down — Executors holding cached blocks stay until cachedExecutorIdleTimeout, which is infinite by default.
- Setting spark.executor.instances as well — It becomes the initial count, which can surprise you.
What you learned
- How Spark adds and removes executors while an application runs
- The configs that control scaling up and down
- Why removing an executor threatens shuffle files, and the three ways around it
- When dynamic allocation helps and when it hurts
Key takeaways
- Dynamic allocation grows executors with the task backlog and releases idle ones.
- Requests ramp up exponentially; idle executors leave after 60 seconds.
- Shuffle files must survive executor removal: external shuffle service, shuffle tracking or decommissioning.
- Set maxExecutors; cached data keeps executors alive.
Check yourself
3 questionsWhy does dynamic allocation need an external shuffle service or shuffle tracking?
Show the answer
Removed executors would take their shuffle files with them. Shuffle output lives on executor disks and is read by later stages.
By default, when is an executor holding cached data released?
Show the answer
Never, unless cachedExecutorIdleTimeout is set. cachedExecutorIdleTimeout defaults to infinity.
Tasks have been pending for a while. How does Spark request executors?
Show the answer
Exponentially: 1, 2, 4, 8 per round. Requests double each round while the backlog persists.
Keep going
Up next · lesson 13 of 30 · 4 min readAdaptive Query Execution
Spark re-plans a running query from real statistics: coalescing, join switching and skew splitting.
Related lessons
Driver, executors and the cluster managerSpark internals · 6 min read
Shuffle internalsSpark internals · 3 min read
Caching and persistence
Previous: Shuffle partitions
Primary sources: Job scheduling: dynamic resource allocation · Configuration: dynamic allocation