Lazy Evaluation: Advanced
Lineage-Based Recovery
A chain whose lineage we can trace
(order_items
.join(products, "product_id")
.filter(F.col("in_stock") == 1)
.groupBy("category")
.agg(F.sum("quantity").alias("units"))
.orderBy(F.col("units").desc()))The Recompute Cost
| Lineage shape | What recovery costs | The risk |
|---|---|---|
| Short and narrow | Replay one short recipe per lost partition | Cheap; fault tolerance is free |
| Very long (iterative) | Replay the entire chain from source | Recovery rivals the original run |
| Wide (post-shuffle) | Re-read and re-shuffle many input partitions | One lost piece pulls in a large fan-in |
Checkpointing to Cut Lineage
| Eager checkpoint | Lazy checkpoint | |
|---|---|---|
| When it writes | Immediately, as a separate action | At the next action that needs the data |
| Extra pass over data | Yes, one dedicated computation | Folded into the next action's run |
| When to prefer | You want the cut materialised now | You want to avoid a separate pass |
The expensive result worth materialising
> From order_items, compute total revenue (quantity times unit_price) per product_id, highest first. The grouping is the wide step whose result you might checkpoint; the aggregation is what makes it worth materialising.
(order_items .("product_id") .agg(F.(F.col("quantity") * F.col("unit_price")).alias("revenue")) .orderBy(F.col("revenue").desc()))
Cache, Checkpoint, or Persist
- Keeps the result in memory/disk
- Avoids recomputing on reuse
- Does NOT cut the lineage
- Evicted or lost blocks recompute from lineage
- Writes the result to reliable storage
- Cuts the lineage behind it
- Survives executor death
- Recovery reads the checkpoint, no replay
Determinism on Replay
| Non-deterministic source | What breaks on replay | The fix |
|---|---|---|
| Random number generation | Different values on recompute | Seed it, or materialise before reuse |
| Current time / timestamps | A different 'now' each replay | Capture the time once, pass it as a value |
| Order-dependent logic | Different result if input order shifts | Make the logic order-independent |
- Treat fault tolerance as cheap only while lineage stays short and narrow.
- Checkpoint to cut a lineage that has grown long or sits past expensive shuffles.
- Use cache/persist for reuse speed and checkpoint for recovery cost; compose them when needed.
- Seed randomness and capture timestamps once so a replay reproduces the original result.
- Don't assume recovery is free; a wide or very long lineage is expensive to replay.
- Don't expect cache to survive executor death; it does not cut the lineage.
- Don't leave non-deterministic transforms in a chain that may be recomputed; replays diverge.
- Don't checkpoint everywhere; it pays a full durable write, so use it where recovery would cost more.
> You own a long iterative Spark job that refines a model over many passes, and it occasionally loses an executor near the end of a multi-hour run. Recovery has been taking almost as long as the original computation, and one run produced numbers that did not reproduce.
A partition is never data Spark trusts to survive. It is a recipe Spark can rebuild.
- Category
- Spark
- Difficulty
- advanced
- Duration
- 15 minutes
- Challenges
- 2 hands-on challenges
Topics covered: Lineage-Based Recovery, The Recompute Cost, Checkpointing to Cut Lineage, Cache, Checkpoint, or Persist, Determinism on Replay
Lesson Sections
- Lineage-Based Recovery (concepts: paSparkExecutionModel)
When an executor dies mid-job, it takes its partitions of in-flight data with it. A naive system would have to start over, because that data is gone. Spark does not, and the reason is that it never treated those partitions as precious irreplaceable data in the first place. It treated them as the output of a known recipe. The lineage is that recipe, recorded per partition, and recovery is just running the recipe again for the partitions that were lost. Say a partition was produced by reading a ch
- The Recompute Cost (concepts: paSparkExecutionModel)
Recovery is cheap when the lineage is short and narrow, because replaying it touches little data and moves none across the network. It gets expensive in two ways, and once you can recognise them you can engineer for the failure instead of trusting fault tolerance blindly. The first is length. Every transformation you chain adds a step to the lineage of the partitions it produces. A pipeline with hundreds of transformations, common in iterative algorithms that loop and refine, builds a very long
- Checkpointing to Cut Lineage (concepts: paSparkCaching)
Checkpointing is the deliberate act of truncating a lineage. When you checkpoint a DataFrame, Spark computes it and writes the result to reliable storage, then discards the lineage behind it. From that point on, the checkpointed data is a new starting point with no history: if a downstream partition is lost, Spark recovers it by reading the checkpoint, not by replaying the entire chain that produced it. You have traded the recompute cost for a one-time write cost. This matters most for the two e
- Cache, Checkpoint, or Persist (concepts: paSparkCaching)
Three operations get confused because they all hold onto a result, though they solve different problems. The confusion is understandable: cache and persist both keep a result around for reuse, and checkpoint also writes a result to storage. What separates them is what they are FOR, and specifically whether they cut the lineage. cache and persist are about speed of reuse. They keep a computed result in memory, or memory and disk, so that the next action reusing it does not recompute the chain, ex
- Determinism on Replay (concepts: paSparkExecutionModel)
Lineage-based recovery rests on an assumption so quiet it is easy to miss, and violating it produces some of the most baffling bugs in Spark: recomputation assumes that replaying a transformation produces the same result it did the first time. If a transformation is non-deterministic, that assumption fails, and recovery can silently produce different data than the partition it is meant to replace. Take a transformation that assigns a random value, or one that depends on the current time, or one