Skip to content
Good engineers know 3 min read

Caching and persistence

cache(), persist() and storage levels: when caching speeds a job up and when it hurts.

You will learn

  • What cache() and persist() actually do, and when
  • The storage levels and their trade-offs
  • When caching speeds a job up and when it slows it down
  • How to release cached data

Read first

Comfortable with these? Read on.

TL;DR cache() marks a DataFrame to be kept in executor memory (spilling to disk) the first time an action computes it, so later actions reuse it. It helps only when the same data is used several times, and it costs memory that joins and aggregations also need.

What it does

Without caching, every action recomputes a DataFrame from its source. df.cache() tells Spark: the next time you compute this, keep the result. It is lazy: nothing is stored until an action runs.

PySparkSpark SQL · Cache once, reuse twice
clean = raw.filter(...).withColumn(...).cache()
clean.count()                 # computes and caches
clean.groupBy("a").count().show()   # reads from cache
clean.write.parquet("out")          # reads from cache
clean.unpersist()
CACHE TABLE clean AS SELECT ... FROM raw WHERE ...;
SELECT a, COUNT(*) FROM clean GROUP BY a;
UNCACHE TABLE clean;

Note the SQL difference: CACHE TABLE is eager by default and computes immediately; CACHE LAZY TABLE waits for the first use.

Storage levels

LevelWhereTrade-off
MEMORY_AND_DISKMemory, spilling to diskDefault for DataFrame.cache(). Safe choice.
MEMORY_ONLYMemory onlyPartitions that do not fit are recomputed when needed.
DISK_ONLYLocal diskCheap on memory; reading back costs I/O.
..._2 variantsTwo copies on two executorsSurvives losing an executor; doubles the space.
OFF_HEAPOff-heap memoryNeeds off-heap memory configured.

DataFrames are cached in Spark's compressed in-memory columnar format, so cached data is often smaller than the same rows as plain objects. Use df.persist(StorageLevel.DISK_ONLY) to choose a level.

When caching helps

  • The DataFrame is used by two or more actions, or two branches of the same job.
  • It is expensive to compute (joins, wide aggregations, slow sources such as JDBC or APIs) and much smaller than its inputs.
  • Interactive exploration, where you query the same subset repeatedly.

When caching hurts

Caching something used once

You pay to write it to memory and gain nothing.

Caching a plain read of Parquet

Reading columnar files is already fast, and caching blocks filter pushdown and column pruning for later queries: they scan the whole cached data instead of skipping files.

Caching huge DataFrames

Storage memory competes with execution memory for joins and aggregations, causing spills, evictions and slower jobs.

Forgetting unpersist

Cached blocks stay until evicted or the application ends. In long-running sessions they accumulate.

Checking the cache

The Storage tab of the Spark UI lists cached DataFrames, the fraction cached and their size in memory and on disk. In the plan, cached data appears as InMemoryRelation / InMemoryTableScan.

Note: checkpoint() is different: it writes the data to reliable storage and cuts the lineage, which helps with very long plans (iterative algorithms). localCheckpoint() does the same on executor storage, faster but not fault tolerant.

Common mistakes

Expecting cache() to compute immediately

It is lazy; the first action fills it.

Caching after the last use

Cache before the DataFrame is reused, not after.

Caching inside a loop without unpersisting

Memory fills with stale versions.

Key takeaways

  • cache() is lazy; the first action stores the data.
  • DataFrame.cache() uses MEMORY_AND_DISK in a compressed columnar format.
  • Cache only data reused across actions and costly to recompute.
  • Unpersist when done; cached data competes with execution memory.

Check yourself

3 questions

1. When is the data of df.cache() actually stored?

Show the answer

At the first action that computes df. cache only marks the DataFrame; the first action fills the cache.

2. Default storage level of DataFrame.cache()?

Show the answer

MEMORY_AND_DISK. DataFrames default to MEMORY_AND_DISK.

3. Why can caching a raw Parquet read slow later filtered queries?

Show the answer

Queries on the cache cannot skip files and row groups using Parquet statistics. Once cached, queries scan the cached data and lose file-level pushdown and pruning.

Go deeper

Primary sources: Caching data in memory · CACHE TABLE