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.
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.