In Spark, what is RDD lineage and how is it used when an executor is lost?
answer
- Spark rebuilds rather than replicates
- every partition knows how it was made
- one parent partition or all of them
- a very long history eventually needs cutting
basics
~20 sLineage is the recorded graph of parent datasets and the deterministic functions that produced each partition. When an executor dies, Spark does not restore a replica: it re-runs that graph to recompute only the lost partitions.
solid answer
~50 sEvery RDD remembers its parents and the function applied to produce each of its partitions; that graph is its **lineage**, and it is Spark's fault-tolerance mechanism. Losing an executor loses its partitions, so the scheduler walks the lineage backwards and recomputes exactly those partitions from whatever still exists. How expensive that is depends on the dependency type: with a **narrow** dependency (`map`, `filter`) each child partition needs one parent partition, so recovery is cheap and local. With a **wide** dependency (a shuffle, from `reduceByKey` or a join) a child partition reads from every map-side output, so if those shuffle files went down with the executor Spark must re-run the whole upstream stage. Caching does not help here — cached blocks are lost with the executor and simply get recomputed. To break a lineage that has grown too long, `RDD.checkpoint()` writes the data to reliable storage and truncates the graph.
code
python · 9 linessc.setCheckpointDir("hdfs:///tmp/spark-checkpoints")
ranks = base_rdd
for i in range(50):
ranks = ranks.join(links).mapValues(update)
if i % 10 == 0:
ranks.cache() # avoid recomputing for the checkpoint write
ranks.checkpoint() # truncate the lineage graph
ranks.count() # materialize so the checkpoint is writtengo deeper
Recall that Spark recovers from failure by recomputing lost partitions from the record of how they were built, rather than by keeping replicas.
Explain the graph itself: parents plus a deterministic per-partition function, and why narrow dependencies recompute cheaply while a lost shuffle output makes the upstream stage re-run.
Demonstrate operating judgment: reading a FetchFailedException correctly, knowing that cache is not durability, running an external shuffle service, and checkpointing iterative jobs before the lineage becomes the failure mode.
Own the reliability tradeoff for the platform: recompute-based recovery is cheap in storage and expensive in time, so decide where to materialize intermediate tables so a failure late in a long pipeline does not cost hours of replay.
## What lineage actually is An RDD is not a dataset sitting somewhere; it is a *recipe*. Each RDD object holds a list of partitions, a reference to its parent RDD (or parents), and the deterministic function that computes a partition from its parents' partitions. Chain three transformations and you have a directed acyclic graph of these recipes — the **lineage graph**, printable with `rdd.toDebugString()`. DataFrames work the same way one level up: the query plan is the lineage, and it compiles down to the same partition-level recipes. This is why Spark can offer fault tolerance without replicating data. A replicated storage system survives failure by having a second copy; Spark survives it by being able to *rebuild* any partition from its inputs, as long as the functions are deterministic and the ultimate source (a file in object storage or HDFS) is still there. ## What happens when an executor dies The driver notices the executor is gone (heartbeat timeout, or the cluster manager reports it). Any tasks in flight there are rescheduled. Any cached blocks it held are marked unavailable. Any shuffle map output files it wrote are marked unavailable too — unless an external shuffle service is serving them, in which case they survive the executor's death. The scheduler then asks: which partitions do I still need, and which of their inputs are missing? It recomputes only what is missing, and only for the partitions actually required. Nothing else in the job is redone. ## Narrow versus wide dependencies decide the cost A **narrow** dependency means each child partition depends on a bounded number of parent partitions — usually exactly one. `map`, `filter`, `mapPartitions`, `union` are narrow. Recovery means recomputing one parent partition chain, and it can often be done on one machine. A **wide** (shuffle) dependency means each child partition may read from *every* parent partition: `reduceByKey`, `groupByKey`, `join` on non-co-partitioned inputs, `repartition`. Here the recovery story changes. Map-side shuffle output is written to local disk, so if it is still readable the reduce side simply refetches it. But if the executor that held those files is gone (and no external shuffle service is running), Spark must re-run the map tasks that produced them — a `FetchFailedException` in the logs is precisely this, and it makes the scheduler resubmit the upstream stage. ## Cache is not a durability mechanism A common misconception: "I cached it, so it is safe". `cache()` / `persist()` is a *performance* hint. Cached blocks live in executor memory (and, with disk levels, on that executor's local disk). Lose the executor and the blocks go with it; Spark quietly recomputes them from lineage. Replicated storage levels such as `MEMORY_AND_DISK_2` keep a second copy on another executor and do reduce recovery cost, at double the space. None of this truncates the lineage — Spark keeps the recipe precisely because the cache can vanish. ## When lineage becomes the problem Iterative workloads — a loop that transforms the same dataset a hundred times, or a graph algorithm — build lineage graphs a hundred levels deep. Two things go wrong. Recovery of a single lost partition may require replaying the entire history. And the driver-side plan or closure graph itself grows large enough to slow down scheduling or blow the stack during serialization. The cure is `RDD.checkpoint()`: set `sc.setCheckpointDir("hdfs://...")` (or any reliable path), call `checkpoint()` on the RDD, and after the next action Spark writes those partitions to that reliable location and **replaces the lineage** with a simple read from it. Recovery afterwards starts from the checkpoint instead of from the source. `localCheckpoint()` truncates lineage using executor storage instead — faster, but it forfeits the durability, so it is only appropriate when a full recompute from source is an acceptable fallback. A common pattern is to `cache()` and then `checkpoint()`, so the checkpoint write does not force a second full recomputation. ## Say "checkpoint" carefully The word means three different things in this space. In core Spark it truncates lineage. In Spark Structured Streaming a checkpoint directory holds offsets and state for a streaming query. In Flink a checkpoint is a periodic distributed snapshot of operator state for exactly-once recovery, and in HDFS checkpointing is the NameNode merging its edit log into a new fsimage. An interviewer will notice if you use the term without saying which system you mean. ## The compact answer "Lineage is the graph of deterministic transformations that produced each partition. It replaces replication as the fault-tolerance mechanism: on executor loss Spark recomputes only the missing partitions. Narrow dependencies make that cheap; a lost shuffle output forces the upstream stage to re-run. Long iterative lineages get truncated with checkpoint, not with cache."
- What is the difference between cache() and checkpoint() for a long-running iterative job?`cache()` keeps the computed partitions in executor memory or local disk but keeps the full lineage, because the cache can be evicted or lost. `checkpoint()` writes the partitions to reliable storage and cuts the lineage above that point, so recovery never replays the earlier history. Cache first, then checkpoint, so the write does not trigger a second full computation.
- You see FetchFailedException in the logs after an executor was killed. What does lineage recovery do?The reduce-side task could not fetch a map output block because the executor that wrote it is gone. The scheduler marks that shuffle output missing and resubmits the *upstream* map stage to regenerate the blocks, then retries the reduce stage. Running an external shuffle service keeps map output readable after an executor dies and avoids this.
- Why must the functions in a lineage be deterministic?Because recovery re-runs them and expects the same output. A closure that reads a mutable global, calls `random` without a fixed seed, or depends on wall-clock time can produce different rows on recompute, silently corrupting results — inconsistent counts between a first run and a retried task are the classic symptom.
Lineage is a receipt showing every step that produced a dish. Lose the dish and you cook it again from the recipe rather than fetching a spare from a fridge that does not exist.
saying these in an interview costs you the question
- Says Spark replicates every RDD partition for fault tolerance
- Thinks cache() makes data durable across executor loss
- Claims a lost partition forces the whole job to restart
- Confuses lineage-truncating checkpoint with a streaming or HDFS checkpoint
- Says wide and narrow dependencies recover at the same cost