Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Lineage and fault tolerance

PySpark · Execution Model

Lineage and fault tolerance

Mediumpyspark-47
lineagefault-tolerancecheckpointshuffle

Question

How does Spark recover when an executor dies mid-job?

Solution

Spark remembers how each partition was built, as a graph called the lineage. If an executor dies and its partitions are lost, Spark reruns just the steps needed to rebuild those partitions, from the source or from the nearest saved data. It does not restart the whole job.

What gets recomputed

Narrow transformations (filter, map, select) are cheap to replay. A lost partition is rebuilt from its single parent partition. The harder case is shuffle output. Map tasks write shuffle files to the executor's local disk. If that executor dies, the files vanish, and reducers cannot fetch them. Spark has to rerun the map tasks that produced the missing files, which may mean going back to an earlier stage.

Stage 0 (map) --writes shuffle files on executor A-->  Stage 1 (reduce)
Executor A dies during stage 1
-> Spark reruns the stage-0 tasks that lived on A, then continues stage 1

An external shuffle service, or on Kubernetes shuffle tracking and decommissioning, keeps shuffle files readable after an executor leaves. That avoids the recompute.

Long lineages

In iterative code (a loop that updates a DataFrame 100 times), the lineage grows long. Recovery then replays a long chain, and even the planning gets slow. checkpoint() writes the data to reliable storage and cuts the lineage there, so recovery starts from the saved copy.

What lineage does not protect

  • The driver. If the driver process dies, the application is lost, because the DAG and scheduler state lived there. Recovery means resubmitting the job, or in streaming, restarting from the checkpoint directory.
  • Side effects. Recomputing a task reruns its code, so a task that wrote to an external database may write twice.
  • A task that fails too many times. After spark.task.maxFailures attempts (4 by default) the job fails.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext