skip to content

Memory Management and Caching

Executor memory is split between execution and storage, and caching decides which of your DataFrames get to occupy it. Most Spark OOM questions are really questions about this split and about what a persist() actually costs.

on this pageshow

explore

questions

6

In Spark, how do the MEMORY_ONLY and MEMORY_AND_DISK persist levels differ when a partition will not fit?

level: juniorimportance: must knowfreq 70%

answer

  1. both cache whole partitions first
  2. the difference is only the fallback
  3. one path re-runs work, the other reads bytes
  4. dropped and recomputed vs written to local disk

basics

~20 s

With MEMORY_ONLY, a partition that does not fit in storage memory is simply not cached and is recomputed from lineage the next time it is needed. With MEMORY_AND_DISK, that partition is written to the executor's local disk and read back instead of recomputed.

solid answer

~40 s

Both levels cache whole partitions in executor storage memory; they differ only in the fallback. `MEMORY_ONLY` drops any partition that does not fit — the next action recomputes it from the lineage of transformations, so a wide, expensive chain can be recomputed repeatedly. `MEMORY_AND_DISK` spills the non-fitting partition to the executor's local disk; the next read deserializes it from there, which is far cheaper than recomputation when the upstream work was heavy. Caching is per-partition and all-or-nothing per partition, so the Spark UI's Storage tab often shows a fraction cached like 63%, not a binary. The defaults differ by API: `RDD.cache()` means `MEMORY_ONLY`, while `Dataset.cache()`/`DataFrame.cache()` means `MEMORY_AND_DISK`. Both levels are lazy — nothing is stored until an action computes the partitions.

code

python · 10 lines
python
from pyspark import StorageLevel

enriched = orders.join(customers, "customer_id").filter("amount > 0")
enriched.persist(StorageLevel.MEMORY_AND_DISK)
enriched.count()          # action: this is when blocks are stored

by_region = enriched.groupBy("region").sum("amount")
by_day = enriched.groupBy("order_date").count()

enriched.unpersist()

go deeper

for a junior

Know that cache() and persist() exist, that they only pay off when a dataset is read more than once, and that the two levels differ in whether a partition that does not fit is recomputed or read back from local disk.

for a middle

Explain the laziness, the per-partition all-or-nothing behaviour, the differing defaults of RDD.cache() versus Dataset.cache(), and what the serialized variants trade CPU for.

for a senior

Show judgment about when not to cache: single-use datasets, caching that blocks filter pushdown, and long-running drivers where forgetting unpersist starves execution memory over hours.

for a principal

Treat caching policy as something a platform decides, not each job: guidance on which intermediates deserve memory, and awareness that uncontrolled caching on shared clusters converts one team's convenience into everyone else's spill.

## What persist actually does `persist(level)` marks a DataFrame or RDD so that the *first* time each of its partitions is computed, the result is kept for reuse. It is **lazy**: the call itself moves no data. Nothing is in memory until an action (`count()`, `write`, `collect()`) forces those partitions to be computed. `cache()` is just `persist()` with the API's default level. Caching happens **per partition**, and per partition it is all-or-nothing: a partition is either fully stored or not stored at all. That is why the Spark UI's *Storage* tab reports a cached fraction — a 200-partition DataFrame can be 63% cached because storage memory filled up part-way through. ## MEMORY_ONLY Partitions are held as deserialized JVM objects (for RDDs) or as Spark's in-memory columnar batches (for DataFrames) in the storage side of the unified memory pool. If a partition does not fit, it is **dropped** — not written anywhere. The next action that needs it re-executes the whole lineage that produced it: re-reading the source files, re-running the filters, and in the worst case re-running an upstream shuffle. This is the level to choose when the upstream work is cheap (a plain file scan) and memory is plentiful; it is the level that quietly turns a "cached" pipeline into one that recomputes on every action. ## MEMORY_AND_DISK Same behaviour while memory lasts; when a partition will not fit, or is later evicted because execution memory reclaimed borrowed space, the block is **written to the executor's local disk directories** (`spark.local.dir` / the cluster manager's scratch dirs) and read back on demand. You trade a disk read for a recomputation. When the upstream chain contains a shuffle, a join, or an expensive UDF, that trade is almost always worth it — which is exactly why the Dataset API defaults to it. Note what disk here is *not*: it is the executor's local scratch space, not durable storage. If the executor dies, the blocks die with it and the partition is recomputed from lineage on a new executor. ## The SER variants `MEMORY_ONLY_SER` and `MEMORY_AND_DISK_SER` keep blocks **serialized** as byte arrays. That typically cuts the footprint substantially and gives the garbage collector one big object instead of millions of small ones, at the cost of CPU on every read. Two caveats worth knowing: - In **PySpark**, storage levels are serialized regardless — Python objects are pickled — so the SER distinction does not apply the way it does in Scala. - For **DataFrames/Datasets**, cached data already uses Spark SQL's compact columnar in-memory format (compressed by default via `spark.sql.inMemoryColumnarStorage.compressed`), so reaching for the `_SER` level buys much less than it does for a raw RDD of case-class objects. ## The rest of the family `DISK_ONLY` skips memory entirely — useful when a result is large, reused a few times, and expensive to compute, but you do not want it competing with execution memory. `OFF_HEAP` stores blocks serialized in the off-heap pool, and requires `spark.memory.offHeap.enabled`. The `_2` suffixed variants (`MEMORY_AND_DISK_2`, etc.) replicate every block on two executors, so losing one executor does not force recomputation. Replication doubles the memory cost and is rarely justified for batch work. ## When caching is the wrong answer Caching earns its cost only when a dataset is read **more than once**. A single linear pipeline that reads, transforms and writes once gains nothing from `cache()` and loses storage memory that execution could have used. Caching also blocks some optimizations: once a plan is cached, downstream filters cannot be pushed through the cache barrier into the source scan, so a `cache()` placed too early can make a job read far more data than the uncached plan would have. ## Releasing the memory `unpersist()` releases the blocks; by default it is non-blocking, and it takes an optional `blocking` flag. Spark's `ContextCleaner` will eventually free cached blocks whose RDD or DataFrame has been garbage-collected on the driver, so short-lived batch applications often get away without calling it. Long-running applications — notebooks, streaming drivers, Thrift servers — must unpersist explicitly, or storage memory slowly fills with datasets nobody references and execution memory starves. ## Reading the evidence The Storage tab shows, per cached dataset, the storage level, the number of cached partitions, the fraction cached, and the split between size in memory and size on disk. A large "size on disk" number next to a `MEMORY_AND_DISK` level tells you memory is the constraint; a cached fraction well under 100% on `MEMORY_ONLY` tells you your job is silently recomputing.

  • Why does the Spark UI show a DataFrame as only 60% cached rather than fully cached or not cached?
    Caching is decided per partition and is all-or-nothing per partition. Once storage memory fills, later partitions are dropped (MEMORY_ONLY) or sent to local disk (MEMORY_AND_DISK) while the earlier ones remain. Executors also evict cached blocks when execution memory reclaims borrowed space, so the fraction can fall during a job. A partially cached dataset still gives correct results — it just recomputes or re-reads the missing partitions.
  • Does calling cache() ever make a Spark job slower?
    Yes. It costs storage memory that execution could otherwise use, so a job that previously fit can start spilling. It also acts as a barrier: filters and column pruning cannot be pushed through a cached plan into the source scan, so caching too early can force reading far more data. And if the dataset is consumed only once, the caching write is pure overhead.
  • When is MEMORY_ONLY_SER worth choosing over MEMORY_ONLY for an RDD?
    When the objects are small and numerous, so per-object JVM overhead dominates and garbage collection pauses are hurting the job. Serialized storage collapses a partition into one byte array, typically cutting footprint and GC pressure at the cost of deserialization CPU on every read. For DataFrames the gain is much smaller, because cached DataFrames already use a compressed columnar in-memory format.

saying these in an interview costs you the question

  • Says MEMORY_ONLY writes leftover partitions to disk
  • Thinks cache() materializes data immediately without an action
  • Claims cache() and persist() use the same default level everywhere
  • Treats the disk in MEMORY_AND_DISK as durable storage that survives the executor
  • Caches every intermediate DataFrame, including single-use ones

context

open as a page

In a Spark executor, what does spark.memory.fraction control and how do execution and storage share it?

level: middleimportance: must knowfreq 62%

basics

~20 s

spark.memory.fraction (default 0.6) sizes the unified pool for execution and storage out of the executor heap left after a fixed 300 MB reservation. Execution and storage borrow from each other on demand; only storage's protected floor is safe from eviction.

open as a page

In the Spark UI, one task reports Spill (Memory) of 48 GB on an executor with an 8 GB heap — what is that measuring?

level: middleimportance: should knowfreq 45%

basics

~20 s

Spill (Memory) is the cumulative in-memory, deserialized size of every buffer the task flushed to disk, summed over all spill events. It is not a peak, so a task that spills repeatedly can report far more than the executor heap ever held.

open as a page

Why does YARN kill a Spark executor for exceeding memory limits when its JVM heap never filled up?

level: seniorimportance: should knowfreq 58%

basics

~20 s

The container limit covers the whole process, not just the heap: spark.executor.memory plus spark.executor.memoryOverhead. Off-heap consumers — shuffle network buffers, JVM metaspace and thread stacks, native libraries and Python workers — burst past the overhead allowance while the heap sits half empty.

open as a page

When several downstream Spark jobs read one expensive DataFrame, how do you choose between persist() and writing it out?

level: principalimportance: should knowfreq 38%

basics

~20 s

Scope decides it. persist() reuses data only within one Spark application and dies with its executors; writing the result to Parquet makes it readable by every later job, survives failures, and keeps executor memory for execution. Cache inside an application, materialize across them.

open as a page

What does setting spark.memory.offHeap.enabled to true change about a Spark executor's memory?

level: seniorimportance: nice to knowfreq 28%

basics

~20 s

It gives Spark a second memory pool of spark.memory.offHeap.size bytes, allocated outside the JVM heap and used by Tungsten for execution and by the OFF_HEAP storage level. That memory is not garbage collected, so large working sets stop driving GC pauses.

open as a page