Skip to content
Good engineers know 4 min read · spark.dynamicAllocation.enabled

Dynamic allocation

Add and remove executors with the workload, and keep shuffle files safe while executors come and go.

TL;DR With 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

  1. Backlog

    Tasks have been pending for schedulerBacklogTimeout (1 s)

  2. Ramp up

    Request 1, then 2, 4, 8 more executors each round while the backlog persists

  3. Cap

    Never beyond what pending tasks need, or maxExecutors

  4. 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.

ConfigDefaultMeaning
spark.dynamicAllocation.enabledfalseTurn it on
...minExecutors / ...maxExecutors0 / unlimitedBounds. Always set a max on shared clusters.
...initialExecutorsminExecutorsExecutors at start (spark.executor.instances is used if larger)
...schedulerBacklogTimeout1sHow long tasks must wait before asking for more
...executorIdleTimeout60sIdle time before an executor is released
...cachedExecutorIdleTimeoutinfinityIdle 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 fitPoor fit
Shared clusters and notebooks with idle time between queriesShort, latency-sensitive jobs: requesting and starting executors takes seconds to minutes
Jobs whose stages vary a lot in task countJobs that cache heavily: cached executors are never released by default
Many applications competing for one clusterLong-running streaming queries with steady load, where a fixed size is simpler and predictable
Note: dynamic allocation scales executors inside one application. Cluster autoscaling (Kubernetes cluster autoscaler, Databricks or EMR autoscaling) adds and removes machines. They work together: Spark asks for executors, and the cluster grows to fit them.

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 questions

Why 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 read
Adaptive Query Execution
Spark re-plans a running query from real statistics: coalescing, join switching and skew splitting.

Related lessons

Previous: Shuffle partitions

Primary sources: Job scheduling: dynamic resource allocation · Configuration: dynamic allocation