A Spark stage fails with FetchFailedException — what happened, and what does Spark do next?
answer
- the error surfaces where the data is missing, not where it was lost
- look at the host named in the stack trace
- the writer's JVM is usually already gone
- recomputing the map stage, not retrying the reader
basics
~20 sA reduce task could not read shuffle blocks from a peer after retries, almost always because the executor holding those files died or the node is unreachable. Spark marks the lost map output as missing, recomputes it, and retries the stage.
solid answer
~40 s`FetchFailedException` is a **read-side symptom of a write-side loss**. A reduce task asked for shuffle blocks at an address from the map output tracker and failed after `spark.shuffle.io.maxRetries` attempts (default 3). Usual causes: the executor was killed (out of memory, or killed by the resource manager for exceeding its memory limit), the node was lost or reclaimed, its local disk filled, a long GC pause blocked the server, or the shuffle service was saturated past `spark.network.timeout`. Spark treats it specially: the `DAGScheduler` unregisters that executor's map outputs, resubmits the parent stage to regenerate only the missing outputs, then retries the reduce stage — it does not count against `spark.task.maxFailures`. After `spark.stage.maxConsecutiveAttempts` (default 4) the job aborts. Fix the executor loss; retrying alone just burns the attempts.
code
text · 6 linesorg.apache.spark.shuffle.FetchFailedException: Failed to connect to host-17:7337
at org.apache.spark.storage.ShuffleBlockFetcherIterator.throwFetchFailedException(...)
Caused by: java.io.IOException: Connection reset by peer
---
# host-17 executor log, minutes earlier:
Container killed by YARN for exceeding memory limits. Consider boosting spark.executor.memoryOverhead.go deeper
Recognise that the error means shuffle data a task needed is unreachable, and that the place to look is the executor named in the message, not the failing stage.
Explain the mechanism: map output is unregistered, the parent stage recomputes the missing partitions, and the retry budget is per stage rather than per task. Name the common causes of executor loss.
Show the diagnosis path end to end — host from the trace, executor exit reason, memory versus congestion versus node loss — and prescribe the structural fix rather than raising timeouts. Interviewers expect you to have debugged this live.
Own the platform answer: shuffle-file durability strategy, executor sizing standards, spot-instance policy and how much recomputation risk the organisation accepts. Frame repeated fetch failures as a capacity and topology decision, not a per-job tuning knob.
## What the exception actually means `FetchFailedException` is raised by the shuffle reader, but the failure lives upstream. The reduce task asked the map output tracker where its partition's blocks are, got an address, and could not get the bytes: connection refused, connection reset, a timeout, or a corrupt/missing block. Spark retried the fetch `spark.shuffle.io.maxRetries` times (default 3, with a wait between attempts) before giving up. The stack trace names the host and port it failed to reach — port 7337 in the message points at the external shuffle service, another port at a peer executor's own block-transfer server. ## Why the map output was gone Walk the causes in the order that finds most incidents: 1. **The executor died.** JVM OOM, a container killed by YARN for exceeding physical memory (the classic "Container killed by YARN for exceeding memory limits" message), or a Kubernetes pod evicted. Without an external shuffle service, the moment the JVM exits its shuffle files become unreachable even though they still exist on disk — nothing is left to serve them. 2. **The node went away.** Spot or preemptible instance reclaimed, node drained, hardware failure. The files really are gone. 3. **Local disk filled.** Shuffle-heavy jobs write a lot under `spark.local.dir`; a full disk fails the write or truncates it, and the reader sees a missing or short block. 4. **The server was alive but unresponsive.** A multi-second stop-the-world GC on the serving executor, or an external shuffle service overwhelmed by thousands of concurrent fetch requests, blows past `spark.network.timeout` (default 120s). 5. **Very large blocks.** A skewed partition can put a huge block behind one fetch; slow transfers and memory pressure on both ends make timeouts likelier. ## What the scheduler does about it This is the part interviewers actually want. A fetch failure is not an ordinary task failure. The `DAGScheduler` receives it and: - unregisters all map outputs held by that executor (and, if the whole host is suspected lost, by every executor on it), so the tracker stops handing out dead addresses; - marks the corresponding **map stage** as having missing partitions and resubmits it — only the missing map tasks rerun, not the whole stage's worth of work if other outputs survive; - resubmits the reduce stage afterwards. Crucially it does **not** count toward `spark.task.maxFailures`, because rerunning the reader would fail identically; the missing input must be regenerated first. What does bound it is `spark.stage.maxConsecutiveAttempts` (default 4): after that many consecutive failed attempts at the stage, the job aborts. That is why a cluster losing executors steadily produces the familiar pattern of a job crawling through repeated stage retries and then dying — each retry recomputes expensive upstream work, so wall-clock cost is superlinear. ## How to diagnose it in practice Do not start in the reduce stage's logs. Take the host named in the exception and look at **that executor's** log and its exit reason. Three readings: - Executor log ends with an OOM or the resource manager's kill message → a memory problem: raise `spark.executor.memoryOverhead` or executor memory, reduce per-task memory demand by increasing shuffle partitions, or lower cores per executor so fewer tasks share the heap. - Executor still alive, failure is a timeout or connection reset → serving-side congestion or GC. Look at the shuffle service logs, block-transfer thread settings and GC pauses on that node. - Whole node absent from the cluster → infrastructure: preemption, autoscaling, hardware. ## How to stop it recurring The durable fixes are structural, not retry-count tuning. Run an **external shuffle service** so files outlive the executor process, or on deployments without one, enable graceful decommissioning so shuffle blocks migrate before an executor goes away. Right-size executors so they are not killed for overhead. Increase shuffle partitions so individual blocks are smaller and no single fetch is enormous. Keep spare local disk. Raising `spark.shuffle.io.maxRetries` and `spark.network.timeout` buys tolerance for transient congestion and is a legitimate mitigation — but it is a band-aid over dying executors, and saying "just increase the retries" is the answer that ends interviews. ## What interviewers are checking That you locate the cause upstream of where the error appeared, that you know the scheduler regenerates map output rather than simply retrying the failed task, and that you can name the bounded retry budget that turns this into a job failure.
- Why doesn't a fetch failure count against spark.task.maxFailures?Because rerunning the reduce task would fail exactly the same way — its input no longer exists. The scheduler must first regenerate the missing map output, so the failure is handled at stage level: outputs are unregistered, the parent stage's missing partitions rerun, then the reduce stage retries. The bound that applies instead is `spark.stage.maxConsecutiveAttempts`, default 4.
- How do you tell an executor OOM from a saturated shuffle service?Follow the host in the exception. If that executor's log ends in an OutOfMemoryError or the resource manager's exceeding-memory kill message, and the executor is gone from the UI, it is memory. If the executor or node is still alive and the failure is a timeout or connection reset, suspect serving-side congestion or a long GC on that node and check the shuffle service logs.
- Does an external shuffle service eliminate fetch failures?No, it removes one large class of them: shuffle files stay fetchable after the executor process exits, so a dead executor no longer invalidates its map output. Losing the whole node, filling the local disk, or overwhelming the service with concurrent fetches still produce fetch failures. It also cannot help if the files were never fully written.
- Why does a job with repeated fetch failures get dramatically slower rather than just slightly?Each fetch failure forces the scheduler to recompute missing map output, which may itself depend on earlier expensive stages. On a cluster losing executors steadily, retries compound: work is redone, new executors take losses of their own, and progress can go backwards until `spark.stage.maxConsecutiveAttempts` aborts the job. Cost grows far faster than the number of failures.
saying these in an interview costs you the question
- Says just increase the retry count and rerun
- Blames the reduce task's code for the failure
- Thinks Spark reruns only the failed reduce task
- Claims shuffle data is replicated so it cannot be lost
- Assumes a fetch failure means a network outage every time