In Spark, how does a task retry differ from a stage retry after FetchFailedException?
answer
- one path is cheap, one is not
- was the input lost, or just the attempt?
- shuffle bytes live on the executor that wrote them
- the scheduler goes back a stage
- find out why the executor died
basics
~20 sA task retry re-runs one failed task on another executor, up to spark.task.maxFailures (default 4). A FetchFailedException means map output is gone, so Spark fails the stage and resubmits the parent stage to regenerate it.
solid answer
~50 sAn ordinary task failure — a bad record, an executor OOM, a lost node — is handled locally: the `TaskSetManager` re-runs **that one task** on another executor, up to `spark.task.maxFailures` attempts (default 4) within the same stage attempt, and only then fails the stage and the job. A `FetchFailedException` is different in kind. It means a reduce-side task could not read shuffle blocks from an executor that has since died or been decommissioned, so the data is not merely unread but **gone**. Retrying that task cannot help. The `DAGScheduler` instead unregisters the lost executor's map output, marks the current stage as failed, and resubmits the **parent** `ShuffleMapStage` to recompute only the missing partitions before re-running the downstream stage. In the Spark UI this appears as a stage with attempt numbers — "retry 1" — and duplicated stage rows. Consecutive stage attempts are capped by `spark.stage.maxConsecutiveAttempts`.
code
text · 5 linesERROR TaskSetManager: Task 118 in stage 7.0 failed:
FetchFailed(BlockManagerId(exec-14, host-9, 7337), shuffleId=3, mapIndex=41, ...)
INFO DAGScheduler: Marking ShuffleMapStage 6 as failed due to a fetch failure
INFO DAGScheduler: Executor lost: 14 -- unregistering its map output
INFO DAGScheduler: Resubmitting ShuffleMapStage 6 and ResultStage 7go deeper
Know that Spark retries a failed task automatically on another executor rather than failing the job immediately, and that the number of attempts is bounded by configuration.
Explain the two mechanisms separately: per-task retry up to spark.task.maxFailures, versus stage resubmission when shuffle output is genuinely lost, and why retrying the reduce task cannot help there.
Trace a real incident: read the fetch failure for the dead executor id, find why it died, and choose between an external shuffle service, partition sizing and instance strategy rather than tuning retry counts.
Own the cost of recomputation at fleet scale — spot-instance policy, shuffle-service or decommission-migration deployment, and the runtime and budget variance that repeated stage retries introduce.
## Two different recovery mechanisms Spark has two independent recovery paths, and confusing them is a common interview stumble. One is per-task and cheap. The other is per-stage and expensive, and it exists because shuffle data is not part of a task's input in the way a source file is. ## Path 1: task-level retry When a task throws — a `NullPointerException` on a malformed record, an executor killed for exceeding its memory limit, a machine that vanished — the `TaskSetManager` for that stage attempt records the failure and re-schedules the same task, preferring a different executor. It repeats up to `spark.task.maxFailures` attempts, default 4. On the fourth failure the stage attempt is aborted and the job fails with the last exception. This is safe because a task's input can be reproduced: for a source stage it re-reads the file or the shuffle input; for a downstream stage it re-fetches shuffle blocks that still exist on disk. Retries are also why Spark tasks must be deterministic and side-effect-free with respect to anything other than Spark's own commit protocol — a task that appends to an external system will append twice if it is retried. Note that failures are counted per task within a stage attempt, and Spark also tracks executors and hosts that repeatedly fail tasks; the blacklisting/exclusion mechanism can stop scheduling onto a consistently bad node so the same task does not burn all four attempts on the same broken disk. ## Path 2: fetch failure and stage retry A shuffle works in two halves: map-side tasks write their output to local disk on the executor that produced it, and reduce-side tasks fetch those blocks over the network. If the producing executor dies after writing, the blocks die with it (unless an external shuffle service or a decommission-migration mechanism preserved them). The reduce task then gets a `FetchFailedException`. No number of retries of the *reduce* task can fix that, because the bytes no longer exist anywhere. So the `DAGScheduler` handles fetch failures specially: 1. It marks the failed reduce stage attempt as failed rather than counting a task failure. 2. It removes the dead executor's registrations from the `MapOutputTracker`, so those map outputs are known to be missing. 3. It resubmits the parent `ShuffleMapStage`, which now has only the missing partitions to recompute — the surviving map output is reused, not regenerated. 4. When the parent completes, it resubmits the downstream stage as a new attempt. That cascade can propagate: if the parent stage's own inputs were themselves shuffle output from a dead executor, the recomputation walks further back up the lineage. Consecutive attempts of the same stage are limited by `spark.stage.maxConsecutiveAttempts`, so a cluster that keeps losing executors eventually fails the job instead of looping forever. ## Reading it in the UI and logs The Stages tab shows the same stage id twice with different attempt numbers, and the second attempt typically has far fewer tasks than the first — those are the partitions that had to be recomputed. In the driver log you see the fetch failure followed by the scheduler explicitly resubmitting stages. The most useful thing to extract is *which executor* the fetch failed from, because the root cause is almost always upstream of Spark: that executor was killed by the resource manager for exceeding its memory, preempted on a spot instance, or lost with its node. ## Root causes worth naming Fetch failures are a symptom, not a disease. The usual causes are: executors killed for exceeding container memory limits (so the fix is memory or partition sizing, not retries), spot/preemptible instance reclamation, node-level disk failures, and network or shuffle-service timeouts under heavy load. Repeated fetch failures that make a job take three times as long are the classic "the job succeeds but the cost has tripled" incident. ## Mitigations Running an **external shuffle service** decouples shuffle output from executor lifetime: blocks are served by a long-lived process on the node, so an executor dying (or being removed by dynamic allocation) no longer destroys its map output. Graceful decommissioning can migrate shuffle blocks off a node that is about to go away. Beyond that, the real fixes address why executors die at all: right-size partitions so tasks do not exceed memory, avoid spot instances for long shuffle-heavy stages, and reduce shuffle volume so there is less to lose. ## The answer in one breath Task retry = re-run one task, bounded by `spark.task.maxFailures` (default 4), used when the input still exists. Stage retry = re-run a whole parent stage to regenerate lost shuffle output, triggered by `FetchFailedException`, bounded by `spark.stage.maxConsecutiveAttempts`, and a signal to go find out why executors are dying.
- What does an external shuffle service change about fetch failures?It moves shuffle block serving out of the executor JVM into a long-lived per-node process, so map output survives an executor exiting — whether it crashed or was removed by dynamic allocation. Reduce tasks keep fetching successfully and no parent stage has to be recomputed. It does not help when the whole node is lost.
- After a fetch failure, why does the retried parent stage usually have far fewer tasks than the first attempt?Because Spark recomputes only the partitions whose map output was lost. The MapOutputTracker still knows about every block on surviving executors, so those are reused. The retry attempt therefore contains just the tasks belonging to the dead executor's share of the partitions, which is why its task count is a fraction of the original.
- Your job succeeds but the Stages tab shows several stages with retry attempts. Is anything wrong?Yes — it succeeded by paying for recomputation. Repeated stage retries usually mean executors are dying, most often killed for exceeding container memory or reclaimed as spot capacity. Track the executor ids named in the fetch failures, check their exit reasons, and fix partition sizing or instance choice rather than accepting the inflated runtime.
saying these in an interview costs you the question
- Says a fetch failure is retried like any other task failure
- Thinks raising spark.task.maxFailures fixes fetch failures
- Believes shuffle files survive on the cluster after an executor dies
- Claims the whole job restarts from the beginning on a stage retry
- Treats FetchFailedException as a network glitch needing no root cause