Lesson 16 / 25

Caching and Persistence

Cache a DataFrame that is reused and release it when finished.

Compute once, reuse many times

Because of lazy evaluation, every action recomputes its whole lineage from the source. If you use the same expensive DataFrame several times (for example in a loop or for multiple outputs), cache() or persist(level) keeps its partitions in executor memory (spilling to disk if needed) after the first action. Cache only data that is reused and expensive to recompute; caching everything wastes memory and can cause evictions. Call unpersist() when done, and remember that caching itself is lazy: it fills on the first action.

Measure, then fix the biggest cost

Most slow jobs are slow because of shuffles, skew, small files or recomputation.

Three levers: cache, join strategy, layout.
Figure 5.1 — Cache, join strategy and layout.

Cache and release, run

I ran this on Apache Spark 4.0.0 (PySpark, local mode, in the official Docker image). After cache() and an action, is_cached is True with the default storage level (memory, spilling to disk). After unpersist() it is False.

c = orders.cache(); c.count()
print(c.is_cached, c.storageLevel)
c.unpersist(); print(c.is_cached)

Output:

True Disk Memory Deserialized 1x Replicated
False

Check the Storage tab

The Spark UI Storage tab shows what is cached and how much memory it uses. If cached blocks keep being evicted, caching is not helping.

Quick check: When is caching worthwhile?

  • When the DataFrame is used once
  • For every DataFrame always
  • When an expensive DataFrame is reused several times
  • Never
Answer

When an expensive DataFrame is reused several times — Caching avoids recomputation, which only helps when the data is reused.