skip to content

Why does spark.dynamicAllocation.enabled require an external shuffle service or shuffle tracking?

level: seniorimportance: should knowfreq 52%

answer

  1. what an executor leaves behind when it exits
  2. the next stage still needs something from it
  3. local disk, not HDFS, not the driver
  4. either the files outlive the process, or the process waits
  5. the Kubernetes answer has a different name

basics

~20 s

Dynamic allocation removes idle executors, but an executor's local disk holds the shuffle files later stages must fetch. Killing it destroys that map output, so Spark needs an external shuffle service to serve the blocks after exit, or shuffle tracking to refuse to remove executors still holding them.

solid answer

~50 s

With `spark.dynamicAllocation.enabled=true`, Spark adds executors while tasks are backed up and removes any that sit idle for `spark.dynamicAllocation.executorIdleTimeout` (60s by default). The catch is that a map task writes its shuffle output to the executor's local disk, and reducers in a later stage fetch it **from that executor's process**. Remove the executor and the blocks become unreachable, producing `FetchFailedException` and expensive stage retries — the job gets slower, not cheaper. Two mechanisms fix it. `spark.shuffle.service.enabled=true` runs an external shuffle service alongside the executors (on YARN, an auxiliary service inside each NodeManager) that keeps serving those files after the executor is gone. Alternatively `spark.dynamicAllocation.shuffleTracking.enabled=true` (Spark 3.0+) makes Spark track which executors still hold shuffle data for active jobs and keep them alive — the usual route on Kubernetes, where there is no shuffle service to install.

code

bash · 8 lines
bash
# -- YARN: external shuffle service serves blocks after an executor exits
spark-submit --master yarn --deploy-mode cluster \
  --conf spark.dynamicAllocation.enabled=true \
  --conf spark.shuffle.service.enabled=true \
  --conf spark.dynamicAllocation.minExecutors=4 \
  --conf spark.dynamicAllocation.maxExecutors=200 \
  --conf spark.dynamicAllocation.initialExecutors=20 \
  rollup.jar

go deeper

for a junior

Know that dynamic allocation lets an application add and drop executors during a run, and that map output lives on executor local disks rather than in shared storage.

for a middle

Explain the add and remove policy with its timeouts, and why removing an executor can strand shuffle blocks that a later stage must fetch. Name both remedies and which platform each belongs to.

for a senior

Diagnose the failure in production: FetchFailedException right after a scale-down, a shuffle service enabled on the driver but missing from the NodeManagers, or an application that never shrinks because cached blocks pin its executors.

for a principal

Own the policy: what maxExecutors a queue permits, whether the platform runs a shuffle service or relies on shuffle tracking, and how decommissioning makes preemptible capacity usable without turning every eviction into a stage recompute.

## What dynamic allocation does By default a Spark application holds the executors it asked for from start to finish, even during a long driver-side plan or a final single-task write. Dynamic allocation (`spark.dynamicAllocation.enabled=true`) lets the application grow and shrink instead: - **Scale up** when there are pending tasks that have been backed up for `spark.dynamicAllocation.schedulerBacklogTimeout` (1s), then repeatedly every `sustainedSchedulerBacklogTimeout`, growing exponentially up to `spark.dynamicAllocation.maxExecutors`. - **Scale down** when an executor has run no tasks for `spark.dynamicAllocation.executorIdleTimeout` (60s), down to `spark.dynamicAllocation.minExecutors` (0 by default). - `initialExecutors` sets the warm start. Setting `spark.executor.instances` (that is, `--num-executors`) at the same time is a configuration conflict — express the ceiling as `maxExecutors`. On a shared cluster this is how a hundred jobs coexist: nobody parks idle containers on a queue while their driver plans, and a backfill can transiently use far more capacity than its steady-state shape. ## Why removing an executor is dangerous Spark's shuffle is a write-then-fetch handoff. In the map stage, each task partitions its output by the reduce key and writes it to files on the executor's **local disk** (`spark.local.dir`) — not to HDFS, not to the driver. The driver's `MapOutputTracker` records where each block lives. In the reduce stage, every reducer opens a connection to every executor that produced a block for it and fetches its slice. That handoff outlives the map task by design: the map stage may finish minutes before the reduce stage begins. If dynamic allocation sees an executor with no running tasks during that gap and hands its container back, the shuffle files go with it. The reducers then fail with `FetchFailedException`, Spark treats it as a lost map output, and it **re-runs the map stage** to regenerate the missing blocks. A feature meant to save resources turns into a retry loop that costs more than the idle executors ever did. ## Fix one: the external shuffle service `spark.shuffle.service.enabled=true` separates serving shuffle blocks from the process that produced them. On YARN, the service is deployed as an auxiliary service inside each NodeManager (the `spark_shuffle` aux-service, with a jar on the NodeManager classpath and matching `yarn-site.xml` entries); on standalone it runs on each worker. Executors write their shuffle files where the service can read them, and reducers fetch from the service, so an executor can exit — voluntarily or not — without stranding its map output. This is the classic pairing on Hadoop clusters, and it has a second benefit: an executor lost to a crash no longer forces a map-stage recompute, because its files are still served. ## Fix two: shuffle tracking Installing a NodeManager aux-service is not possible on Kubernetes, so Spark 3.0 added `spark.dynamicAllocation.shuffleTracking.enabled=true`. Instead of making the blocks survive the executor, Spark refuses to remove an executor that still holds shuffle data referenced by an active job, with `spark.dynamicAllocation.shuffleTracking.timeout` as a backstop for data nothing will read again. It is simpler to operate and needs nothing installed on the nodes, but it is weaker: executors holding shuffle output stay alive even when idle, so the scale-down you get is less aggressive than with a shuffle service. Spark 3.1 and later can also decommission executors gracefully — migrating shuffle and cached blocks off a node before it goes away (`spark.decommission.enabled` and the related storage-decommission settings) — which is what makes dynamic allocation on spot or preemptible instances tolerable. ## Cached data has the same problem An executor holding cached RDD or DataFrame blocks is in the same position: evict it and the cache is gone, forcing recomputation. Spark treats this separately — `spark.dynamicAllocation.cachedExecutorIdleTimeout` defaults to infinity, so executors holding cached blocks are not removed for idleness at all. If you cache aggressively and then wonder why the application never scales down, that is why. ## Operating it - Set `minExecutors` above zero for latency-sensitive or streaming work; scaling from zero adds container-start latency to the first stage. - Set `maxExecutors` deliberately — an unbounded ceiling on a shared queue is how one backfill starves everything else. - Watch the executor add/remove timeline in the Spark UI (or the History Server) when a job looks slow: a sawtooth of executors being added and removed within a minute usually means the timeouts fight the job's stage rhythm. - If you see `FetchFailedException` shortly after executors are removed, check that the shuffle service is actually enabled on the nodes — a mismatch between the driver's `spark.shuffle.service.enabled` and what the NodeManagers run is a real failure mode.

  • Why do executors holding cached blocks often refuse to be scaled down?
    They are governed by `spark.dynamicAllocation.cachedExecutorIdleTimeout`, which defaults to infinity so cached data is never thrown away for mere idleness. Removing them would force recomputation of everything the user asked Spark to cache. If you want them released, set that timeout explicitly or unpersist when the cached data is no longer needed.
  • What are sensible values for minExecutors and maxExecutors on a shared cluster?
    Set `minExecutors` above zero when start-up latency matters — streaming jobs and interactive sessions should not scale to nothing between micro-batches. Set `maxExecutors` to a real ceiling reflecting the queue share the job is entitled to; leaving it unbounded lets one backfill consume the queue and starve everyone behind it.
  • How does dynamic allocation interact with --num-executors?
    They express the same thing two ways and should not both be set: `--num-executors` fixes `spark.executor.instances`, while dynamic allocation wants a range. Use `spark.dynamicAllocation.initialExecutors` for the warm start and `minExecutors`/`maxExecutors` for the bounds; leaving a stale `--num-executors` in a submit script is a common source of confusion about why scaling looks wrong.

Shuffle files are like parcels left in a courier's van. Dynamic allocation sends idle vans home; the external shuffle service is a depot that keeps the parcels once the van leaves, while shuffle tracking simply refuses to send a van home while parcels are still aboard.

saying these in an interview costs you the question

  • Thinks shuffle files are written to HDFS or the driver
  • Says removing an idle executor is always free
  • Believes dynamic allocation works out of the box on Kubernetes
  • Confuses the shuffle service with the block manager cache
  • Sets both --num-executors and dynamic allocation bounds

context