skip to content

For a shuffle-heavy Spark platform on elastic infrastructure, where should shuffle files live?

level: principalimportance: nice to knowfreq 30%

answer

  1. a stateless engine that is quietly stateful
  2. how disposable is your compute?
  3. the bytes outlive the process, or they do not
  4. decouple storage from compute — at a price

basics

~20 s

Choose by how disposable your compute is. Executor-local disk is fastest but dies with the executor; a per-node external shuffle service outlives executors; block migration on decommission covers planned removal; a disaggregated shuffle service decouples shuffle storage from compute entirely.

solid answer

~50 s

Spark shuffle output is written to executor-local disk and is unreplicated, so its durability is whatever the compute's durability is — and on preemptible or aggressively autoscaled infrastructure that is very little. Four options, in increasing decoupling: **local disk only**, fastest and cheapest but a lost executor forces recomputation of its map output; **an external shuffle service** per node (`spark.shuffle.service.enabled`), which keeps files fetchable after an executor exits and is the classic enabler for scaling executors down; **graceful decommissioning**, which migrates shuffle blocks off an executor before it goes away, covering planned removal but not sudden loss; and a **disaggregated shuffle service** such as Apache Celeborn or Apache Uniffle, which pushes shuffle data off the compute node entirely. Decide on preemption rate, job size, recomputation cost and how much operational surface you are willing to run.

code

properties · 6 lines
properties
# node-local service: files survive executor exit
spark.shuffle.service.enabled=true

# planned removal: move blocks before the executor goes
spark.decommission.enabled=true
spark.storage.decommission.shuffleBlocks.enabled=true

go deeper

for a junior

Know that shuffle files are written to the local disk of the executor that produced them and are not replicated, so losing that executor means the data must be recomputed.

for a middle

Explain why an external shuffle service lets files outlive the executor process, and why that is what makes removing idle executors safe during a job.

for a senior

Diagnose recomputation cascades on preemptible infrastructure and pick the mitigation the deployment actually supports — shuffle service, decommission migration, or better sizing — with evidence from retried stage time.

for a principal

Own the tradeoff across the fleet: preemption rate, job duration, shuffle volume, tenant isolation and the operational bill of running a separate shuffle tier. Sequence the change and state what measurement would justify each step.

## Why this is a platform decision, not a job setting Every Spark shuffle materialises map output to the local disk of the executor that produced it, with no replication. That design is fast and simple, and it silently assumes the compute stays alive long enough for the readers to arrive. On fixed long-lived clusters that assumption held. On spot instances, aggressive autoscaling and short-lived pods it does not, and the consequence is repeated fetch failures with recomputation cascades that make long jobs slower the more the platform scales down. Where shuffle data lives is therefore an infrastructure decision that determines how much elasticity you can actually exploit. ## Option 1 — executor-local disk only The default. Lowest latency, no extra components, and every byte stays on the machine that wrote it until fetched. The cost is coupling: when the JVM exits, its shuffle files stop being served even though the bytes are still on disk, and when the node vanishes they are genuinely gone. Fine when nodes are stable and jobs are short enough that recomputation after an occasional loss is cheap. Provisioning implication: shuffle-heavy work wants fast local NVMe and enough of it, which pushes you toward specific instance shapes. ## Option 2 — an external shuffle service per node A long-lived process on each node serves shuffle files out of the executors' local directories, so files outlive the executor process. On YARN it runs inside the NodeManager as an auxiliary service; on standalone deployments it runs on each worker. This is the mechanism that historically made scaling executors down safe: you can remove an idle executor without stranding the shuffle output it already produced. It does not save you from losing the whole node, and the service itself becomes a shared bottleneck — thousands of concurrent fetches against one process is a real congestion source, and a saturated service produces exactly the timeouts that look like network problems. ## Option 3 — decommissioning and block migration When an executor is about to be removed and you have warning — a scale-down decision, a spot-termination notice — Spark can migrate shuffle blocks off it before it exits, and dynamic allocation can alternatively be told to track shuffle ownership and simply not remove executors that still hold needed files. This covers **planned** departure well and costs no extra infrastructure. It cannot help with sudden loss, and migration takes time proportional to the data, so on a heavily preempted fleet the notice window may be shorter than the copy. ## Option 4 — a disaggregated shuffle service Projects such as Apache Celeborn and Apache Uniffle move shuffle data off the compute node into a separate service, typically with push-style writes that aggregate a partition's data before the readers arrive. The compute becomes genuinely stateless: an executor or node can disappear at any moment without invalidating anything, so spot instances and rapid scale-down stop being risky. It also converts the reduce side's many-small-random-reads pattern into larger sequential ones, which is a real win at very high partition counts. The costs are honest ones: another distributed system to size, monitor, upgrade and capacity-plan; another network hop for every shuffle byte; and storage cost for the shuffle tier. Related in spirit, push-based shuffle merges map output into per-partition files on the shuffle service so readers fetch fewer, larger blocks. ## How to actually decide Drive it from measurements, not preference: - **Preemption or scale-down rate.** If executors routinely disappear mid-job, coupling shuffle to compute is already costing you; quantify the recomputation in retried stage-hours. - **Job shape.** A fleet of ten-minute jobs can absorb losses cheaply. Multi-hour jobs with terabyte shuffles cannot: one late fetch failure can rerun hours of work. - **Shuffle volume per job.** A disaggregated tier is easier to justify when shuffle bytes dominate the workload; on light shuffles the extra hop is pure overhead. - **Multi-tenancy.** A shared per-node service is a shared failure and contention domain — one abusive job's fetch storm degrades neighbours. Disaggregation gives you a place to enforce isolation and quotas. - **Operational capacity.** Running another stateful distributed service is a standing cost. If the team cannot operate it well, a well-provisioned local-disk fleet with correctly sized executors beats a poorly run shuffle tier. ## The staged answer A credible plan sequences rather than jumps: first fix executor sizing and partition counts so executors stop dying of self-inflicted memory pressure; then enable whichever executor-outliving mechanism your deployment supports; then, if measured recomputation from node loss is still material and shuffle volume justifies it, evaluate a disaggregated service on the largest workloads before adopting it fleet-wide. ## What interviewers are checking Whether you can reason about state in a supposedly stateless compute layer, whether you attach the decision to preemption rates and recomputation cost rather than to novelty, and whether you acknowledge the operational bill of the most decoupled option.

  • Why does losing executors hurt long shuffle-heavy jobs disproportionately?
    Missing map output must be regenerated, and that regeneration may depend on earlier expensive stages, so one fetch failure can rerun hours of work. On a fleet that keeps reclaiming nodes, retries compound and the job can lose ground faster than it gains it, eventually exhausting the stage attempt budget and aborting after consuming enormous compute.
  • What does a disaggregated shuffle service cost you, honestly?
    Another stateful distributed system to size, monitor, upgrade and page on; an extra network hop for every shuffle byte; and storage capacity for the shuffle tier. It also adds a new failure domain shared across tenants. It pays off when preemption-driven recomputation and many-small-block fetch overhead are both large and measured, not by default.
  • Where would you start before changing shuffle storage architecture at all?
    With the self-inflicted losses. Executors killed for exceeding memory, partition counts left at the default, and skewed stages produce most fetch failures on typical platforms, and no shuffle-storage choice fixes them. Correct sizing and partitioning first, measure the residual loss rate caused purely by infrastructure reclaiming nodes, then decide whether decoupling is warranted.

saying these in an interview costs you the question

  • Assumes Spark replicates shuffle data for durability
  • Treats the compute layer as stateless during a shuffle
  • Adopts a remote shuffle service without measuring recomputation cost
  • Ignores that a shared shuffle service is a contention domain
  • Believes scaling executors down is always safe by default

context