What you will be able to do
- Tell cache() from persist() and state the default storage level for DataFrames and for RDDs
- Choose a storage level from the memory-versus-CPU trade-off
- Explain how cached data is evicted, freed and shared across sessions
- Relate the dataset cache and shuffle spills to executor memory pressure and garbage collection
1.cache() versus persist() and their defaults
Executors are long-lived JVM processes that can keep data in memory or on disk between tasks. Spark lets you use this to keep a dataset around across operations. When you persist a dataset, each node stores the partitions it computes and reuses them in later actions on that dataset or on datasets derived from it. The RDD guide says this often makes later actions more than 10x faster. Caching is a key tool for iterative algorithms and interactive work.
cache() takes no arguments and always uses the default storage level. persist() accepts a StorageLevel, and called with no argument it uses the same default. The default depends on which API you are using:
- DataFrame cache() and persist() default to MEMORY_AND_DISK_DESER. The PySpark default changed to this in 3.0 to match Scala.
- RDD cache() is shorthand for StorageLevel.MEMORY_ONLY, which stores deserialized objects in memory.
You can assign a storage level with persist() only if the DataFrame does not already have one.
Neither call does any work immediately. The level takes effect the first time the DataFrame is computed. The data is kept from that first action onward and reused afterwards.
df = spark.range(1)
df.persist()
# DataFrame[id: bigint]
df.explain()
# == Physical Plan ==
# InMemoryTableScan ...
from pyspark.storagelevel import StorageLevel
df.persist(StorageLevel.DISK_ONLY)
# DataFrame[id: bigint]Checkpoint 1 of 5· Fill the gap
Complete the call that keeps this DataFrame's data on disk only.
from pyspark.storagelevel import StorageLevel
df.persist(StorageLevel. ? )The reference example passes StorageLevel.DISK_ONLY to persist() to store the data on disk.
Source: docs.databricks.comCheckpoint 2 of 5· Exam question
A shared Databricks notebook already has an active `SparkSession` when a second cell runs the following code: ```python spark2 = SparkSession.builder.appName("secondary-analysis").getOrCreate() print(spark2 is spark) ``` What does this print, and why?
Correct answer: A — `True`, because `getOrCreate()` detects the existing active session for the JVM and returns that same object instead of constructing a new one, regardless of the `appName` passed.
- A. This is correct behavior: `getOrCreate()` checks for an active session on the current thread/JVM and returns it if one exists, so `spark2 is spark` evaluates to `True` and the `appName` argument on the second call is ignored.
- B. The `appName` argument only applies when a new session is actually being built; it has no effect on an already-active session and does not trigger a teardown-and-recreate cycle.
- C. There is no duplication of the underlying `SparkContext`: `getOrCreate()` returns the identical object, so only one `SparkContext` and its memory footprint exist, not two.
- D. Databricks notebook cells in the same attached cluster session share the same driver JVM, so they can and do reference the same `SparkSession` object across cells.
- E. Calling `getOrCreate()` multiple times is the documented, safe way to obtain the current session; it does not raise an exception on repeated calls within the same JVM.
2.Storage levels: trading memory against CPU
A storage level decides where persisted partitions live, whether they are serialized and whether they are replicated. The RDD programming guide lists the full set. MEMORY_AND_DISK_DESER, the DataFrame default, is named in the DataFrame reference, but that table does not include it.
| Storage level | Behaviour |
|---|---|
| MEMORY_ONLY | Deserialized Java objects in the JVM; partitions that don't fit are not cached and are recomputed when needed. Default for RDDs. |
| MEMORY_AND_DISK | Deserialized Java objects in the JVM; partitions that don't fit are stored on disk and read from there |
| MEMORY_ONLY_SER (Java and Scala) | Serialized Java objects, one byte array per partition: more space-efficient, more CPU-intensive to read |
| MEMORY_AND_DISK_SER (Java and Scala) | Like MEMORY_ONLY_SER, but spills partitions that don't fit to disk instead of recomputing them |
| DISK_ONLY | Partitions stored only on disk |
| MEMORY_ONLY_2, MEMORY_AND_DISK_2, etc. | Same as the levels above, but each partition is replicated on two cluster nodes |
| OFF_HEAP (experimental) | Like MEMORY_ONLY_SER, but in off-heap memory, which must be enabled |
In Python the serialized/deserialized distinction does not apply, because Python always serializes stored objects with Pickle. The Python storage levels are MEMORY_ONLY, MEMORY_ONLY_2, MEMORY_AND_DISK, MEMORY_AND_DISK_2, DISK_ONLY, DISK_ONLY_2 and DISK_ONLY_3.
To choose a level, the guide gives a sequence:
1. If the data fits comfortably at MEMORY_ONLY, stay there. It is the most CPU-efficient option.
2. If it does not fit, try a serialized level with a fast serializer (Java and Scala).
3. Spill to disk only when computing the partitions was expensive or filtered out a lot of data. Otherwise recomputing a partition can be as fast as reading it from disk.
4. Use a replicated _2 level when you need fast fault recovery.
Every level is fault-tolerant, because Spark can recompute lost partitions from the transformations that created them. Replication only saves you the wait for that recomputation.
Checkpoint 3 of 5· Check yourself
A cached dataset is cheap to recompute and slightly too large for memory. Following the guide, which choice is least justified?
The guide advises against spilling to disk unless computing the dataset was expensive or filtered out a lot of data. For cheap data, recomputing can be as fast as reading from disk.
“Don’t spill to disk unless the functions that computed your datasets are expensive, or they filter a large amount of the data.”Source: spark.apache.org
3.How long cached data lives, and who can see it
Cached partitions are not permanent. Spark monitors cache usage on each node and drops old partitions in least-recently-used (LRU) order. To free a dataset yourself, call unpersist(). By default it does not block. If you need to wait until the resources are actually freed, pass blocking=true.
The cache is also scoped more widely than one session: cached data is shared across all Spark sessions on the cluster. Cached data is still tied to the executors that hold it. When a worker is decommissioned, for example under autoscaling, its Spark cache is lost, and Spark must reread the missing partitions from the source.
| Feature | Disk cache | Apache Spark cache |
|---|---|---|
| Stored as | Local files on a worker node | In-memory blocks, but it depends on storage level |
| Applied to | Any Parquet table stored on S3, ABFS, and other file systems | Any DataFrame or RDD |
| Triggered | Automatically, on the first read (if cache is enabled) | Manually, requires code changes |
| Evaluated | Lazily | Lazily |
| Evicted | Automatically in LRU fashion or on any file change, manually when restarting a cluster | Automatically in LRU fashion, manually with unpersist |
Checkpoint 4 of 5· Match them up
Match each Spark-cache behaviour to its description
Tap a term, then the definition that fits it.
The Spark cache is evicted by LRU or unpersist(), shared across sessions and stored on the workers, so it disappears when they go away.
“Spark automatically monitors cache usage on each node and drops out old data partitions in a least-recently-used (LRU) fashion.”Source: spark.apache.org
4.Executor memory and garbage collection
Memory used for caching and memory used for computation come from the same place: the executor JVM. Databricks describes the executor as a Java process that triggers GC lazily, and notes that the dataset cache uses a lot of executor memory. Large caches therefore reduce the memory left for everything else, which is one reason to unpersist() data you no longer need instead of waiting for LRU eviction.
Shuffles add pressure of their own. Operations such as reduceByKey and aggregateByKey build in-memory structures to organize records. When the data does not fit, Spark spills those tables to disk, which adds disk I/O and increases garbage collection. Shuffles also write intermediate files that are kept until the corresponding RDDs are no longer used and have been garbage collected. If the application keeps references to those RDDs, or GC runs infrequently, cleanup can be delayed for a long time. Long-running jobs can then use a lot of disk space in the directory set by spark.local.dir.
Checkpoint 5 of 5· Check yourself
A long-running Spark job keeps consuming local disk space even though old shuffle stages finished long ago. What is the most likely reason according to the docs?
Shuffle files are kept until their RDDs are no longer used and have been garbage collected. Retained references or infrequent GC delay that, so disk use grows.
“if the application retains references to these RDDs or if GC does not kick in frequently”Source: spark.apache.org
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.DataFrame.cache() stores data at MEMORY_ONLY, the same as RDD.cache().Why is that wrong?
MEMORY_ONLY is the RDD default. For DataFrames, cache() and persist() default to MEMORY_AND_DISK_DESER.
Covered in cache() versus persist() and their defaults
2.You can call persist() again with a different StorageLevel to change how an already-persisted DataFrame is stored.Why is that wrong?
persist() can assign a new storage level only when the DataFrame does not have one yet.
Covered in cache() versus persist() and their defaults
3.unpersist() frees the cached memory before it returns.Why is that wrong?
unpersist() does not block by default. Pass blocking=true to wait until the resources are freed.
Practise it for real
See a DataFrame's persistence show up in its physical plan
1.In a PySpark session, run df = spark.range(1)
Why: range() creates a small DataFrame with a single id column to experiment on
You should see: A DataFrame with schema [id: bigint]
2.Call df.persist() with no argument
Why: With no level given, persist() uses the DataFrame default, MEMORY_AND_DISK_DESER
You should see: DataFrame[id: bigint] is returned
3.Run df.explain()
Why: The physical plan shows whether the persisted data will be read
You should see: The plan contains InMemoryTableScan
Stuck? Get a nudge
If InMemoryTableScan is missing, check that you called explain() on the same df variable you persisted.
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“If no storage level is specified defaults to (MEMORY_AND_DISK_DESER).”
↩︎ cache() versus persist() and their defaults“Sets the storage level to persist the contents of the DataFrame across operations after the first time it is computed.”
↩︎ cache() versus persist() and their defaults“Storage level to set for persistence. Default is MEMORY_AND_DISK_DESER.”
↩︎ Storage levels: trading memory against CPU“Cached data is shared across all Spark sessions on the cluster.”
↩︎ How long cached data lives, and who can see it“Databricks recommends moving away from DataFrame.persist() as it is not compatible with Databricks serverless compute architecture.”
↩︎ How long cached data lives, and who can see it“This can only be used to assign a new storage level if the DataFrame does not have a storage level set yet.”
↩︎ Exam trap 2 - 2.
“The default storage level has changed to MEMORY_AND_DISK_DESER to match Scala in 3.0.”
↩︎ cache() versus persist() and their defaults“Persists the DataFrame with the default storage level (MEMORY_AND_DISK_DESER).”
↩︎ Exam trap 1“Persists the DataFrame with the default storage level (MEMORY_AND_DISK_DESER).”
↩︎ Prediction - 3.https://spark.apache.org/docs/latest/rdd-programming-guide.htmlSecondary source
“The cache() method is a shorthand for using the default storage level, which is StorageLevel.MEMORY_ONLY (store deserialized objects in memory).”
↩︎ cache() versus persist() and their defaults“Caching is a key tool for iterative algorithms and fast interactive use.”
↩︎ cache() versus persist() and their defaults“In Python, stored objects will always be serialized with the Pickle library, so it does not matter whether you choose a serialized level.”
↩︎ Storage levels: trading memory against CPU“Spark’s storage levels are meant to provide different trade-offs between memory usage and CPU efficiency.”
↩︎ Storage levels: trading memory against CPU“the replicated ones let you continue running tasks on the RDD without waiting to recompute a lost partition”
↩︎ Storage levels: trading memory against CPU“Spark will spill these tables to disk, incurring the additional overhead of disk I/O and increased garbage collection”
↩︎ Executor memory and garbage collection“these files are preserved until the corresponding RDDs are no longer used and are garbage collected”
↩︎ Executor memory and garbage collection“This means that long-running Spark jobs may consume a large amount of disk space.”
↩︎ Executor memory and garbage collection“Note that this method does not block by default. To block until resources are freed, specify blocking=true when calling this method.”
↩︎ Exam trap 3“Don’t spill to disk unless the functions that computed your datasets are expensive, or they filter a large amount of the data.”
↩︎ Checkpoint“Spark automatically monitors cache usage on each node and drops out old data partitions in a least-recently-used (LRU) fashion.”
↩︎ Checkpoint“if the application retains references to these RDDs or if GC does not kick in frequently”
↩︎ Checkpoint - 4.
“When a worker is decommissioned, the Spark cache stored on that worker is lost.”
↩︎ How long cached data lives, and who can see it“The Databricks disk cache differs from Apache Spark caching.”
↩︎ How long cached data lives, and who can see it - 5.
“The Apache Spark executor is a Java process that triggers GC lazily”
↩︎ Executor memory and garbage collection“the Apache Spark dataset cache uses a lot of Apache Spark executor memory”
↩︎ Executor memory and garbage collection