Spark Fault Tolerance via RDD Lineage
When a Spark executor fails, the partitions it held are lost. Instead of replicating data (like HDFS), Spark uses lineage: the recorded sequence of transformations that produced each partition.
The driver detects the failure, identifies the lost partitions, and recomputes them by replaying the transformations from the nearest persisted ancestor. This approach is cheaper than replication for most workloads because transformations are typically narrow (map, filter) and operate on local data.
A reduceByKey job on three workers: one worker's executor or whole node fails, Spark traces the lineage back to the shuffle files on disk, and a new executor recomputes only the partitions that were lost.
Click a worker to fail it, and choose whether only its executor crashes or the whole node is lost. An executor crash takes the partitions cached in its memory, but its shuffle files stay on the node’s disk and serve as the recovery point; a lost node takes those shuffle files too, and only the map tasks that wrote them are replayed from HDFS.