Dynamic allocation lets Spark add executors when tasks are waiting and release them when they sit idle, so an application does not hold a fixed number of executors for its whole life.
How it behaves
You turn it on with spark.dynamicAllocation.enabled=true and set bounds:
spark.dynamicAllocation.minExecutors = 2 spark.dynamicAllocation.maxExecutors = 50 spark.dynamicAllocation.executorIdleTimeout = 60s
If tasks are pending, Spark requests more executors, ramping up in steps (1, 2, 4, 8 ...) until it hits the maximum. When an executor has been idle past the timeout, Spark gives it back to the cluster.
The shuffle files problem
An executor that is removed would normally take its shuffle files with it, and later stages would need them. So dynamic allocation needs one of two things: an external shuffle service running on each node (the older way on YARN), or shuffle tracking (spark.dynamicAllocation.shuffleTracking.enabled), which keeps executors that still hold needed shuffle data alive. Shuffle tracking is what you use on Kubernetes, where there is no external service by default.
When it is a good fit
- Shared clusters, where idle executors waste capacity other teams could use.
- Bursty workloads, such as a notebook that runs a heavy query and then sits quiet for twenty minutes.
- Pipelines with stages of very different size: a wide stage with 2,000 tasks followed by a small write.
When it is not
Short jobs. If a job runs for 90 seconds, the ramp-up takes a significant part of that. Many teams fix the size for short, predictable jobs, and use dynamic allocation for long or interactive ones. It also makes cost less predictable, since the executor count varies. Always set a maximum, or a mistake in one job can take over the whole cluster.