Calling cache() in Spark is one line, and most of what decides whether it helps happens out of sight. A cached DataFrame is not one object in one place. It is hundreds or thousands of blocks spread across executors, each tracked by the driver, each subject to eviction and each lost if its executor goes away. Understanding that layer explains most real cache surprises: jobs that slow down after an autoscaling event, a cache that is "mostly" there, and a cache that was never as big as anyone thought.

This article is about that layer. Choosing what to cache, the storage levels, plan matching and the choice between persisting and checkpointing are covered in Spark persistence, and the memory regions the cache lives in are covered in Spark memory management. Here the focus is the machinery: blocks, the BlockManager, the driver's bookkeeping, failure and scaling, and how to size a cache from measurements.

Advertisement

cache() and persist() are the same call

cache() is persist() with the default storage level, and both only mark the data; nothing is stored until an action computes partitions. The default differs by API. RDD.cache() uses MEMORY_ONLY. Dataset.cache() uses MEMORY_AND_DISK (PySpark names its DataFrame default MEMORY_AND_DISK_DESER). SQL's CACHE TABLE registers a named entry and is eager unless written CACHE LAZY TABLE.

For a DataFrame, persisting registers the query plan with the session's cache manager. Later queries whose plan contains an equivalent subtree are rewritten to read the cache instead of recomputing it. For an RDD, persisting sets a storage level on that RDD object, and only computations going through that exact object benefit.

From partition to block

When a task computes a partition of a persisted RDD, it asks its executor's BlockManager for the block first. Blocks of cached RDDs are named rdd_<rddId>_<partitionIndex>, which is exactly what you see in the Storage tab and in executor logs. On a hit the task reads the block, from local memory, local disk or a remote executor. On a miss it computes the partition from its parents and offers the result to the BlockManager to store at the requested level.

A DataFrame cache is an RDD underneath. The cached relation's plan is executed once into an RDD of cached batches, and that RDD is persisted, so the same block machinery applies; the block ids belong to that internal RDD.

Where a cached partition lives, and who knows about itDriverCacheManagerplan to cache entryBlockManagerMasterblock to executorsExecutor 1MemoryStorerdd_12_0, rdd_12_3DiskStorerdd_12_7BlockManagerreports blocks, serves remote readsExecutor 2MemoryStorerdd_12_1 (lost if executor dies)block status updatesschedule task for partition 3 hereremote fetchMissing block?recompute from lineageLineage is the fallback for any lost or evicted blockReplicated levels (_2) keep a second copy on another executor
Blocks are stored by each executor's BlockManager and reported to the driver, which schedules later tasks where their blocks are. Anything missing, evicted or lost is recomputed from lineage.
Advertisement

Memory store, disk store and unrolling

Each BlockManager has a memory store and a disk store. A block stored deserialized is held as Java objects, the fastest form to read and the largest. A serialized level (MEMORY_ONLY_SER on RDDs) stores bytes, which are smaller but must be deserialized on every read. A disk level writes the serialized bytes to the executor's local directories.

Storing a partition in memory is not instant, because Spark does not know its size until it has iterated it. The memory store unrolls the iterator incrementally, reserving unroll memory as it goes, starting with spark.storage.unrollMemoryThreshold (1 MiB by default). If it cannot reserve enough, the block does not fit. Under a memory-only level the partition is then simply not cached and will be recomputed next time; under MEMORY_AND_DISK it is written to disk instead. This is why a cache can be partially present: some partitions fitted, others did not, and the Storage tab shows a fraction cached below 100 percent.

The cache lives in the storage part of unified memory: spark.memory.fraction (0.6) of the heap after a 300 MB reserve is shared between execution and storage, and spark.memory.storageFraction (0.5) of that is protected from eviction by execution.

The columnar cached batch

DataFrames are not cached as rows. Spark SQL builds columnar batches of up to spark.sql.inMemoryColumnarStorage.batchSize rows (10,000 by default) and, with spark.sql.inMemoryColumnarStorage.compressed (true by default), chooses a compression scheme per column from statistics, such as run-length or dictionary encoding. Each batch also carries per-column statistics, minimum and maximum, so a filter over the cache can skip whole batches.

Two consequences matter. A cached DataFrame is usually much smaller than the same data as Java objects, often comparable to the Parquet size, so estimates based on row objects overshoot. And a scan of the cache reads only the columns the query needs. Larger batches compress better but need more memory to build; if caching a very wide table runs out of memory while building batches, lowering the batch size is the first lever.

How the driver tracks the cache

Every BlockManager reports block status changes to the driver's BlockManagerMaster: which block, which executor, which store and how large. The scheduler uses those locations as locality preferences, so a task for partition 3 prefers the executor that holds rdd_12_3 in memory, waiting up to spark.locality.wait before running elsewhere and fetching the block remotely or recomputing it.

The driver is also what removes caches. unpersist() tells every executor to drop the blocks and is non-blocking by default; pass blocking=True when the next step needs the memory immediately. With spark.cleaner.referenceTracking enabled (the default), an RDD that becomes unreachable in the driver program is eventually cleaned up by the context cleaner after garbage collection, but relying on that for large caches makes memory release timing depend on driver GC. Unpersist explicitly.

# PySpark: build, cache, materialize, verify, then release
enriched = (spark.read.parquet("s3://lake/events/dt=2026-10-01")
            .join(spark.table("dim_users"), "user_id")
            .select("user_id", "country", "plan", "event_type", "ts"))

enriched.persist()                 # Dataset default: MEMORY_AND_DISK
rows = enriched.count()            # full action: every partition computed and stored
print(enriched.storageLevel)       # confirm the level actually set

daily = enriched.groupBy("country").count()
daily.explain()                    # expect InMemoryTableScan, not a FileScan + join

by_plan = enriched.groupBy("plan", "event_type").count().collect()

enriched.unpersist(blocking=True)  # free storage memory before the next heavy stage

Executor loss and replication

Blocks in an executor's memory or local disk are lost when that executor dies, whether from an out-of-memory kill, a lost spot instance or a node failure. The cache metadata on the driver is updated, and the next task needing a lost partition recomputes it from lineage, at the cost of rereading sources and redoing any shuffle above it. Nothing fails; the job just slows, often sharply, which is why the symptom is a stage that was fast on the first run and slow on the third.

Replicated storage levels such as MEMORY_AND_DISK_2 store a second copy on another executor. With spark.storage.replication.proactive (true by default), Spark replenishes lost replicas from a surviving copy. Replication doubles the memory cost, so reserve it for caches that are expensive to rebuild and running on unreliable capacity such as spot instances.

Dynamic allocation and decommissioning

Dynamic allocation releases idle executors, and an executor holding cached blocks is idle in every sense except that it holds your cache. Spark therefore uses a separate timeout for executors with cached data, spark.dynamicAllocation.cachedExecutorIdleTimeout, which defaults to infinity. The effect is that a single forgotten cache can pin executors for the life of the application. Either unpersist when done or set a finite timeout and accept recomputation.

With spark.shuffle.service.fetch.rdd.enabled (false by default) and the external shuffle service, disk-persisted RDD blocks can be served by the shuffle service, so executors holding only disk blocks can be released after the normal idle timeout. Graceful decommissioning goes further: with spark.decommission.enabled and spark.storage.decommission.enabled, a decommissioning executor migrates its blocks to other executors before exiting; RDD block migration is controlled by spark.storage.decommission.rddBlocks.enabled, which defaults to true. See Spark dynamic allocation for the scaling side.

Reading the Storage tab

The Storage tab of the Spark UI is the authoritative view of what is cached. Each entry shows its storage level, the number of cached partitions against the total, the fraction cached, and size in memory and on disk. Clicking an entry lists every block with its executor and location, which is where you see that partition 41 lives on disk on executor 7 while the rest are in memory. The Executors tab shows storage memory used against available per executor, which reveals a cache concentrated on a few executors because their tasks happened to run first.

Read the numbers in this order. If the cached partition count is below the total, either an action did not touch every partition (a take() or limit() computes only what it needs) or partitions failed to fit. If size on disk is non-zero under MEMORY_AND_DISK, memory was short and reads of those partitions pay disk and deserialization cost. If an entry you expect is missing, the persist was never followed by an action, or it was unpersisted, or a different DataFrame object was cached.

Worked example: sizing a cache

A job joins 120 GB of Parquet events with a user dimension and reuses the result in four aggregations. The cluster has 20 executors with 16 GB heaps. Storage plus execution memory per executor is (16 GB minus 300 MB) times 0.6, about 9.4 GB, so about 188 GB across the cluster, of which half, 94 GB, is protected for storage.

Rather than guess, measure. Cache a 5 percent sample, materialize it with count(), and read Size in Memory from the Storage tab. Suppose it shows 3.1 GB for 6 GB of Parquet, so the full cache would be about 62 GB in memory: it fits in the protected region with room for execution. Had it shown 140 GB, the options would be caching only the columns the four aggregations use, filtering first, using a disk level, or writing an intermediate Parquet table instead. After the full run, check that fraction cached reads 100 percent and that the aggregations show InMemoryTableScan in their plans.

Failure modes

  • Partial cache: fraction cached below 100 percent after a memory-only persist. Partitions that did not unroll are recomputed every time. Use a disk-backed level or cache less.
  • Cache never used: no InMemoryTableScan in the plan because the query was rebuilt with a different plan. Reuse the same DataFrame object.
  • Slowdown after scaling or spot loss: blocks lost with executors and recomputed from lineage. Replicate expensive caches or enable decommissioning.
  • Executors that never scale down: cached blocks with an infinite cached-executor timeout. Unpersist or set a finite timeout.
  • Execution spills while caches sit idle: large caches occupying storage memory that execution cannot reclaim below the protected fraction. Unpersist with blocking before heavy stages.
  • Stale results: the cache is a snapshot of its inputs at computation time; refresh or uncache after the sources change.

Trade-offs

Caching trades memory and the risk of loss for skipping recomputation. It pays when the cached data is reused several times and is expensive to rebuild, such as the result of a wide join or shuffle, and it hurts when the cache is used once, crowds out execution memory or hides a better plan. For data that must survive executor loss or be reused across applications, writing a table or checkpointing is more robust than any storage level; lineage itself is explained in Spark RDD lineage.

What to do next

  1. List every persist call in your jobs and confirm each one is reused at least twice and unpersisted when no longer needed.
  2. Materialize caches with a full action and check the Storage tab for 100 percent fraction cached and the expected size.
  3. Check plans with explain() for InMemoryTableScan where you expect cache hits.
  4. Size large caches from a measured sample, not from row counts, and compare with the protected storage memory.
  5. If you use dynamic allocation, set spark.dynamicAllocation.cachedExecutorIdleTimeout deliberately and consider decommissioning for spot capacity.
  6. Use a replicated level only for caches that are expensive to recompute and run on unreliable executors.
Key takeaway: A Spark cache is a set of named blocks held by executor BlockManagers and tracked by the driver, not a single object. Blocks can fail to fit, be evicted or vanish with an executor, and Spark quietly recomputes them from lineage. Materialize caches fully, measure their real size, verify plans use them, unpersist them deliberately, and decide in advance how a cache should behave under executor loss and autoscaling.