Spark does not keep results around. Every action, whether a count(), a write or a collect(), turns the logical plan into stages and runs them from their inputs. If three outputs are derived from the same expensive join, Spark runs that join three times unless you ask it to keep the result. That is what persistence is for: persist() and its shorthand cache() tell Spark to store the partitions of a dataset the first time they are computed and to serve later reads from that copy.

Caching the wrong thing wastes memory that shuffles and joins needed, a cache that was never fully materialized silently recomputes, and forgotten caches pin executors that dynamic allocation would otherwise release. This article explains what happens underneath, how to choose a storage level, how to size a cache, what fails, and how to decide between caching, checkpointing and simply writing a table. The memory model itself is covered in Spark memory management and recomputation in RDD lineage; here we focus on the persistence decision.

Advertisement

Why recomputation is the default

Spark's fault tolerance rests on lineage: every dataset knows the transformations that produce it from stable inputs, so a lost partition can be rebuilt by replaying them. Because of that, Spark never needs to store intermediate results for correctness, and by default it does not. Transformations are lazy and build a plan; actions execute it from the leaves.

The one exception is shuffle output. Spark automatically keeps map-side shuffle files on local disk so that a failed reduce task does not force the whole upstream stage to rerun. The Spark documentation still recommends calling persist on a result you plan to reuse, because shuffle files cover only the stage boundary, not the work after it. For how shuffle files are written and served, see Spark shuffle.

What persist() actually does

Calling persist(level) does no work. It records a storage level on the dataset and returns. For the RDD API the level is attached to the RDD object. For DataFrames and Datasets the driver's CacheManager registers the query's analyzed plan together with an InMemoryRelation placeholder; any later query whose plan matches that plan is rewritten to read from the cache.

The data is stored the first time a task computes a partition. The task asks its executor's BlockManager for the block, named like rdd_42_7 for partition 7 of RDD 42. On a miss it computes the partition, hands the result to the BlockManager, and the BlockManager keeps it in the memory store, the disk store, or both, according to the level. For DataFrames the stored form is not rows of Java objects but compressed columnar batches, controlled by spark.sql.inMemoryColumnarStorage.compressed (true by default) and spark.sql.inMemoryColumnarStorage.batchSize (10,000 rows per batch by default).

Blocks live on the executor that computed them, and later tasks prefer that executor. A cache dies with its application.

Where a persisted partition lives, and what happens when it is missingDriverCacheManager + DAGAction runs a taskneeds partition 7BlockManager lookuprdd_42_7 on this executor?Memory storestorage region of heapDisk storelocal dirs, *_AND_DISKRemote replicaonly with _2 levelsRecomputereplay lineageSource / shufflefiles, shuffle outputscheduleget blockhitspill or readmissno replicaread inputsA miss is never an error: Spark falls back to lineage. The cost of a miss is the whole upstream stage.Eviction under memory pressure drops MEMORY_ONLY blocks and moves MEMORY_AND_DISK blocks to disk.
The read path for a persisted partition. A cache miss is silent and costs a full recomputation of the lineage behind that partition.
Advertisement

Storage levels and the defaults that differ by API

A storage level combines four choices: use memory, use disk, keep data deserialized as objects or serialized as bytes, and how many replicas to store. The levels documented in the Spark programming guide are:

LevelWhat it storesWhen memory runs out
MEMORY_ONLYDeserialized objects in the JVM heapPartitions that do not fit are not cached and are recomputed on each use
MEMORY_AND_DISKDeserialized objects, overflow on local diskPartitions that do not fit are written to disk and read back
MEMORY_ONLY_SER, MEMORY_AND_DISK_SEROne serialized byte array per partition (Java and Scala only)As above; smaller, more CPU to read
DISK_ONLYLocal disk onlyNot applicable
*_2, DISK_ONLY_3Same as the base level with two (or three) replicas on different nodesLost executor does not force recomputation
OFF_HEAPSerialized, off-heap memory (experimental)Requires off-heap memory to be enabled

The defaults are the first trap, because they are different in each API. RDD.cache() means MEMORY_ONLY. Dataset.cache() in Scala and Java means MEMORY_AND_DISK. In PySpark, DataFrame.cache() means MEMORY_AND_DISK_DESER, which matches the Scala behaviour. PySpark RDDs are always pickled, so the serialized and deserialized distinction does not exist there, and only MEMORY_ONLY, MEMORY_ONLY_2, MEMORY_AND_DISK, MEMORY_AND_DISK_2, DISK_ONLY, DISK_ONLY_2 and DISK_ONLY_3 are offered for them.

The programming guide's own guidance is a good starting rule. Keep the default if it fits. If it does not, use a serialized level with a fast serializer such as Kryo (see Kryo serialization). Spill to disk only when the computation behind the data is expensive or filters out a lot, because otherwise recomputing a partition can be as fast as reading it back from disk. Use replicated levels only when you cannot afford to wait for recomputation after an executor loss.

A worked example: one join, three outputs

A daily job reads about 400 GB of event parquet, keeps three event types, joins a small user dimension and produces three marts: hourly counts by country, a session funnel and revenue by segment. Without persistence, each of the three writes triggers its own job, and each job scans 400 GB and repeats the join. With the enriched DataFrame persisted, the scan and join run once.

from pyspark import StorageLevel
from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder.appName("sessions-daily").getOrCreate()

events = spark.read.parquet("s3://lake/events/dt=2026-09-29/")      # ~400 GB of parquet
users = spark.read.parquet("s3://lake/dim/users/")

# The expensive, reused part: filter, join, derive. Cache AFTER the filter and projection.
enriched = (events
    .where(F.col("event_type").isin("view", "add_to_cart", "purchase"))
    .select("user_id", "session_id", "event_type", "ts", "sku", "price")
    .join(F.broadcast(users.select("user_id", "country", "segment")), "user_id")
    .withColumn("hour", F.date_trunc("hour", "ts")))

enriched.persist(StorageLevel.MEMORY_AND_DISK_DESER)   # PySpark's DataFrame default, stated explicitly
n = enriched.count()          # full action: materializes every partition (show() would not)
print("rows cached:", n, "is_cached:", enriched.is_cached)

by_hour    = enriched.groupBy("hour", "country").count()
funnel     = enriched.groupBy("session_id").agg(F.collect_set("event_type").alias("steps"))
revenue    = enriched.where("event_type = 'purchase'").groupBy("segment").agg(F.sum("price"))

for name, df in [("by_hour", by_hour), ("funnel", funnel), ("revenue", revenue)]:
    df.write.mode("overwrite").parquet(f"s3://lake/marts/{name}/dt=2026-09-29/")

enriched.unpersist()          # non-blocking by default; pass blocking=True in tests

Three details matter. The persist call comes after the filter and the projection, so the cache holds six columns of the rows that are used, not the raw table. The count() is a deliberate full action: show() or take(10) only compute the partitions they need, so they leave the cache mostly empty and the first real job pays for the rest. And the unpersist() at the end releases executor memory for whatever runs next in the same application.

Now size it. Say the executors have 16 GiB of heap each and there are 20 of them. With the default spark.memory.fraction of 0.6 and 300 MB reserved, each executor has roughly (16,384 MB - 300 MB) x 0.6, about 9.4 GiB, of unified memory shared by execution and storage, or about 188 GiB across the cluster. By default half of it, about 4.7 GiB per executor, is the part of storage that execution cannot evict (spark.memory.storageFraction = 0.5). If the Storage tab of the Spark UI shows the enriched DataFrame at 120 GiB in memory, it fits, but it takes most of the unified pool, and the three aggregations will have to share what is left for their shuffles. If it shows 250 GiB, the MEMORY_AND_DISK family will spill part of it to local disk, which is fine for this job because the alternative is a 400 GB scan plus a join.

Read the real size from the Storage tab after materialization; parquet size on disk is a poor predictor.

Materialization and plan matching

Because persist is lazy, the first action pays for computing and storing. That catches people in two ways. First, as above, partial actions produce a partial cache. Second, for DataFrames the cache is found by plan matching. A query uses the cache only if its plan contains a subtree equivalent to the one registered. Reusing the same DataFrame variable is safe. Rebuilding the same logic in another function with a different column order, a different literal or a re-read of the source can produce a plan that does not match, and it silently recomputes.

Check with explain(). A cache hit appears as InMemoryTableScan over an InMemoryRelation. If you see the original file scan and join, the cache is not being used. The SQL surface behaves slightly differently: CACHE TABLE is eager by default and scans immediately, while CACHE LAZY TABLE behaves like persist.

-- SQL surface: CACHE TABLE is eager unless you say LAZY
CACHE TABLE enriched_today OPTIONS ('storageLevel' 'MEMORY_AND_DISK')
  AS SELECT * FROM enriched_view;
CACHE LAZY TABLE dim_users;
UNCACHE TABLE IF EXISTS enriched_today;

-- In the plan, a cache hit shows up as InMemoryTableScan over an InMemoryRelation
EXPLAIN SELECT country, count(*) FROM enriched_today GROUP BY country;

A cache is also a snapshot. It holds the data as it was when each partition was computed. If the source files change, the cached DataFrame keeps serving the old rows until you uncache it or refresh the table with REFRESH TABLE or spark.catalog.refreshTable. In long-running applications this is a correctness issue.

Caching can also make a query slower. A direct read can prune partitions and push filters into the parquet reader; a query over a cached relation scans the cached batches instead. Cache what is reused as it is, not a superset every query has to filter again.

Eviction and how the cache competes with execution

Storage and execution share the unified memory region. Storage can grow into free execution memory, and execution can take memory back by evicting cached blocks, but only down to the protected storage fraction. Blocks are evicted in least-recently-used order, and a block from the same RDD being written is not evicted to make room for itself. What eviction means depends on the level: a MEMORY_ONLY block is simply dropped and will be recomputed on the next access, while a MEMORY_AND_DISK block is written to the disk store first.

The consequence is that caches and shuffles compete. A large cache that sits inside the protected storage fraction reduces what a sort or hash aggregation can use before it spills. If shuffle spill grows after adding a cache, shrink the cache: narrow the projection, serialize it or use DISK_ONLY. The full picture of the regions, and of how off-heap memory fits in, is in Spark memory management.

Failure modes

  • Executor loss. Cached blocks live in executor memory and local disk, so when an executor dies its blocks are gone. The next access recomputes them from lineage, which can mean rerunning the most expensive stage in the job. Replicated levels avoid the wait at twice the memory cost.
  • Dynamic allocation that never scales down. Executors that hold cached blocks are not released by default, because spark.dynamicAllocation.cachedExecutorIdleTimeout defaults to infinity. A forgotten cache keeps a whole fleet alive. Unpersist explicitly or set a finite timeout.
  • Disk-persisted blocks after executor exit. With spark.shuffle.service.fetch.rdd.enabled and the external shuffle service, blocks persisted to disk can be served after their executor is removed; see the external shuffle service. On graceful decommissioning, spark.storage.decommission.rddBlocks.enabled migrates cached blocks to other executors.
  • Caches in loops. Iterative code that persists a new DataFrame each iteration without unpersisting the previous one fills storage with dead blocks, and each iteration's lineage grows. Unpersist the previous generation and checkpoint every few iterations.
  • Caching before the filter. Persisting the raw source and filtering afterwards stores far more than is used and can evict the blocks that matter.

Persist, checkpoint or write a table?

These three are often confused. persist keeps data on executors for the life of the application and keeps the lineage, so it is fast and fault tolerant by recomputation, but lost on executor loss and on restart. checkpoint() writes the data to the reliable directory set by setCheckpointDir and truncates the lineage, which is what you want when the plan has grown so long that planning or recomputation becomes the bottleneck; localCheckpoint() truncates lineage but stores on executors, so it is fast and not fault tolerant. Writing an explicit table in parquet, Delta or Iceberg is the right answer when another job, another team or tomorrow's run needs the result.

Rule of thumb: persist for reuse within one job, checkpoint to cut a long lineage, write a table to share across jobs.

Operational guidance

Make caches visible. The Storage tab of the Spark UI lists each persisted dataset with its level, the fraction of partitions cached, and its size in memory and on disk. A fraction below 100 percent after a full action means the level could not hold it and partitions are being recomputed or read from disk. df.is_cached and df.storageLevel in PySpark, and spark.catalog.isCached for tables, let you assert the state in code and tests.

Record cached size per job and review unpersist calls as you would review closing a file.

What to do next

  1. List the DataFrames your job uses more than once; persist only those, after filters and projections.
  2. Set the storage level explicitly instead of relying on the API default, and use MEMORY_AND_DISK for anything expensive to recompute.
  3. Materialize with a full action such as count() and confirm 100 percent cached in the Storage tab.
  4. Run explain() on each downstream query and look for InMemoryTableScan.
  5. Compare shuffle spill before and after caching; shrink or move the cache to disk if spill grows.
  6. Unpersist when the last consumer finishes, and set a finite cachedExecutorIdleTimeout if you use dynamic allocation.
  7. Replace caches that must survive failures or be shared with a checkpoint or a written table.
Key takeaway: Persistence turns a recomputed plan into stored blocks on executors. It is lazy, it is found by plan matching, it competes with execution for the same memory and it disappears when an executor or the application does. Cache only datasets that are reused, after filtering and projection, with an explicit storage level. Materialize with a full action, verify with the Storage tab and explain(), and unpersist when done. When the result must survive failures or be shared, checkpoint or write a table instead.