Learn how Spark stores intermediate DataFrames and RDDs in memory, disk, or both—and when persistence can dramatically improve performance.
If you work with Apache Spark long enough, you’ll eventually encounter a situation like this:
You have a DataFrame that takes several minutes to compute, and you use it multiple times in your pipeline.
Without persistence, Spark may recompute that DataFrame every time it is needed.
That’s where persistence comes in.
Persistence tells Spark to save the result of a computation so that it can be reused instead of being recomputed.
But persistence is not simply “put everything in memory.”
Spark provides multiple storage levels, allowing you to control whether data is stored in memory, on disk, serialized, or replicated.
In this article, we’ll understand:
- What persistence means in Spark
- Why persistence is needed
persist()vscache()- Spark storage levels
- How persistence works internally
- Persistence with DataFrames
- Persistence with RDDs
- When you should use persistence
- When persistence can hurt performance
- How to remove persisted data
- Production best practices
- Common interview questions
📚 Apache Spark Learning Path
Previous: [Caching in Spark]
Current: Persistence in Spark
Next: [Checkpointing in Spark]
Complete: The Complete Apache Spark Learning Path: Beginner to Advanced
1. What Is Persistence in Spark?
Spark uses lazy evaluation.
When you write:
df = spark.read.parquet("sales/")df = df.filter(df.amount > 1000)df = df.groupBy("region").sum("amount")
Spark doesn’t immediately execute these transformations.
Instead, it builds a logical execution plan.
The actual computation happens when an action such as:
df.show()
or
df.count()
is executed.
Now imagine that the resulting DataFrame is used multiple times:
df.count()df.groupBy("region").count().show()df.groupBy("product").sum("amount").show()
Without persistence, Spark may need to execute the upstream computation again when different actions require the same intermediate data.
Persistence allows us to tell Spark:
“Compute this dataset and keep the result around because I will use it again.”
2. Why Do We Need Persistence?
Consider this pipeline:
sales = ( spark.read.parquet("sales/") .filter("amount > 1000") .join(customers, "customer_id"))
Suppose sales is used by several downstream operations:
sales.groupBy("region").sum("amount").show()sales.groupBy("product").sum("amount").show()sales.groupBy("customer_segment").count().show()
The expensive operations might include:
Read Parquet ↓Filter ↓Join ↓sales DataFrame ↙ ↓ ↘region product segment
If Spark repeatedly has to execute the expensive upstream portion, the same work may be performed multiple times.
We can persist the intermediate result:
sales.persist()
Then trigger an action:
sales.count()
After the computation completes, Spark stores the generated partitions according to the selected storage level.
Subsequent operations can reuse those persisted partitions.
3. persist() vs cache()
This is one of the most common Spark interview questions.
Both methods are used to keep computed data for reuse.
cache()
df.cache()
is essentially shorthand for persistence using Spark’s default storage level.
For DataFrames, the default storage level is:
MEMORY_AND_DISK_DESER
depending on the Spark API/version semantics.
persist()
persist() gives you explicit control over the storage level.
For example:
df.persist(StorageLevel.MEMORY_ONLY)
or:
df.persist(StorageLevel.MEMORY_AND_DISK)
So the key difference is:
| Method | Purpose |
|---|---|
cache() | Use the default storage level |
persist() | Choose a specific storage level |
Think of it this way:
cache() ↓"I want Spark to keep this dataset."persist() ↓"I want Spark to keep this dataset using this specific storage strategy."
4. Spark Storage Levels
Spark provides different storage levels.
The major ones you’ll encounter are:
MEMORY_ONLYMEMORY_ONLY_2MEMORY_AND_DISKMEMORY_AND_DISK_2DISK_ONLYDISK_ONLY_2OFF_HEAP
The exact available options can vary slightly by Spark API and version.
Let’s understand the important ones.
5. MEMORY_ONLY
With:
from pyspark import StorageLeveldf.persist(StorageLevel.MEMORY_ONLY)
Spark attempts to store the dataset in memory.
If there isn’t enough memory to store all partitions, partitions that don’t fit are not automatically written to disk.
They can therefore be recomputed when needed.
Conceptually:
Dataset
|
-------------------
| | |
Partition Partition Partition
↓ ↓ ↓
RAM RAM RAM
Advantage
Very fast when the entire dataset fits comfortably in memory.
Disadvantage
If memory is insufficient, some partitions may have to be recomputed.
This is an important distinction.
MEMORY_ONLY does not mean:
“Store everything in RAM no matter what.”
It means:
“Store partitions in memory when possible.”
6. MEMORY_AND_DISK
This is one of the most useful storage levels.
df.persist(StorageLevel.MEMORY_AND_DISK)
Spark stores partitions in memory when possible.
If a partition doesn’t fit into memory, Spark spills that partition to disk.
Conceptually:
Dataset
|
-------------------
| | |
↓ ↓ ↓
RAM RAM Disk
When the dataset is too large to fit completely in memory, this can be safer than MEMORY_ONLY.
The tradeoff is obvious:
Memory → FasterDisk → Slower
But disk persistence can still be much faster than repeatedly recomputing an expensive transformation pipeline.
7. DISK_ONLY
With:
df.persist(StorageLevel.DISK_ONLY)
Spark stores the persisted partitions on disk.
This is useful when:
- the dataset is too large for memory
- recomputation is expensive
- memory is needed for other workloads
However, disk access is slower than memory access.
Therefore, you should not automatically choose DISK_ONLY.
8. Serialized vs Deserialized Storage
Another important concept is whether Spark stores cached data in:
Deserialized form
or:
Serialized form
Serialized data generally requires less memory but may require CPU overhead for serialization/deserialization.
Conceptually:
DeserializedObject → Object → Object ↓ Larger memory Faster access
versus:
SerializedObject → Bytes ↓ Smaller memory CPU needed to deserialize
This creates a classic tradeoff:
Memory efficiency ↔ CPU overhead
Modern Spark’s DataFrame execution uses its own optimized internal representations, so you should not blindly reason about DataFrame persistence exactly like old-style Python RDD caching.
9. Persistence with DataFrames
Let’s look at a practical example.
Suppose we have:
sales = spark.read.parquet("sales/")
We perform an expensive transformation:
processed_sales = ( sales .filter("amount > 1000") .join(customers, "customer_id"))
Now suppose we need this DataFrame multiple times.
We can persist it:
from pyspark import StorageLevelprocessed_sales.persist(StorageLevel.MEMORY_AND_DISK)
Then materialize it:
processed_sales.count()
Now Spark has actually computed the DataFrame and populated its persisted partitions.
We can reuse it:
processed_sales.groupBy("region").sum("amount").show()processed_sales.groupBy("product").sum("amount").show()processed_sales.groupBy("customer_segment").count().show()
10. Why Do We Need an Action After persist()?
This is another important concept.
Consider:
df.persist()
You might assume Spark immediately computes and stores the DataFrame.
It doesn’t.
persist() is associated with the DataFrame/RDD and tells Spark how to store its computed partitions when they are materialized.
Because Spark uses lazy evaluation, you still need an action:
df.count()
or:
df.show()
or:
df.write.parquet(...)
The first action computes the required partitions and populates the persistence layer.
Conceptually:
df.persist() | ↓Mark dataset for persistence | ↓Action | ↓Compute dataset | ↓Store partitions | ↓Reuse persisted data
11. Persistence Does Not Break the Lineage
This distinction is extremely important.
Suppose we have:
df1 = spark.read.parquet("sales/")df2 = df1.filter("amount > 1000")df2.persist()
Persistence stores computed partitions for reuse.
But Spark still maintains the underlying lineage.
That means persistence is primarily a performance optimization.
It is not the same thing as checkpointing.
This gives us an important distinction:
| Persistence | Checkpointing |
|---|---|
| Mainly performance optimization | Mainly lineage/truncation and fault-recovery mechanism |
| Stores computed partitions | Writes data to checkpoint storage |
| Lineage remains | Checkpoint can truncate lineage |
| Often used for repeated reuse | Useful for very long/complex lineage |
| Usually executor storage | Checkpoint storage |
We’ll explore checkpointing in the next article.
12. Persistence with RDDs
Persistence was originally a fundamental concept in Spark’s RDD model.
Example:
rdd = sc.textFile("sales.txt")processed = ( rdd .map(parse_record) .filter(is_valid))
We can persist it:
processed.persist(StorageLevel.MEMORY_ONLY)
Then:
processed.count()
Once materialized, Spark can reuse the stored partitions.
RDDs also allow explicit storage-level selection.
For example:
processed.persist(StorageLevel.MEMORY_AND_DISK)
13. Persistence Happens at the Partition Level
Spark doesn’t generally think of a DataFrame as one giant object stored somewhere.
Spark distributes data into partitions.
For example:
DataFrame | +-- Partition 0 +-- Partition 1 +-- Partition 2 +-- Partition 3
When you persist a dataset, Spark stores its computed partitions according to the selected storage level.
For example:
Executor 1 ├── Partition 0 → Memory └── Partition 1 → MemoryExecutor 2 ├── Partition 2 → Memory └── Partition 3 → Disk
This is why executor memory and partition sizing matter when working with persistence.
14. Persistence and Executors
Spark executors are responsible for executing tasks.
They also provide storage space for persisted data.
Conceptually:
Spark Application
|
---------------------
| |
Executor 1 Executor 2
| |
Cached data Cached data
This means persistence consumes executor resources.
If you persist a huge dataset, you are effectively asking Spark to reserve resources for that dataset.
That leads to an important production rule:
Never persist data simply because you can. Persist data because reuse justifies the storage cost.
15. When Should You Use persist()?
Persistence is most useful when a dataset is:
1. Expensive to compute
For example:
Read large dataset ↓Multiple filters ↓Join ↓Aggregation ↓Complex transformation
If this result is reused several times, persistence can eliminate repeated computation.
2. Used by multiple downstream operations
For example:
processed.persist()processed.groupBy("region").count()processed.groupBy("product").count()processed.groupBy("customer").count()
This is a strong persistence candidate.
3. Used repeatedly in iterative algorithms
Machine learning workflows can repeatedly process the same data.
For example:
Iteration 1 → DatasetIteration 2 → DatasetIteration 3 → Dataset...Iteration N → Dataset
Persistence can avoid rebuilding the same intermediate dataset repeatedly.
4. Shared across multiple actions
For example:
processed.count()processed.write.parquet(...)processed.groupBy(...).count()
If the upstream computation is expensive, persistence can help.
16. When Should You NOT Use Persistence?
This is just as important.
Case 1: Dataset is used only once
For example:
df = spark.read.parquet("sales/")df.filter("amount > 1000").write.parquet("output/")
Persisting the DataFrame probably adds unnecessary overhead.
Case 2: Dataset is tiny
If the computation is extremely cheap, persistence may provide little benefit.
The overhead of storing and managing the cached data may not be worthwhile.
Case 3: Dataset is too large
If you persist a dataset that consumes most of your executor memory, you can create memory pressure.
That can cause:
Memory pressure ↓Evictions ↓Garbage collection ↓Spills ↓Slower jobs
Case 4: You persist everything
This is one of the most common Spark mistakes.
Bad pattern:
df1.persist()df2.persist()df3.persist()df4.persist()df5.persist()
Persistence is not free.
Every persisted dataset consumes resources.
17. Always Unpersist When You’re Done
Once the persisted dataset is no longer needed:
df.unpersist()
Example:
processed.persist()processed.count()processed.groupBy("region").count().show()processed.groupBy("product").count().show()processed.unpersist()
This tells Spark that the persisted data can be removed.
For long-running applications, this is particularly important.
18. unpersist() and Memory Management
Suppose your application does:
Stage 1 ↓Persist dataset A ↓Use A ↓Stage 2 ↓Persist dataset B ↓Use B
If dataset A is no longer needed but remains persisted:
Executor memory----------------------Dataset A █████████Dataset B █████████Other data █████----------------------
You may unnecessarily reduce the memory available for active computations.
A better pattern is:
df_a.persist()# use df_adf_a.unpersist()df_b.persist()# use df_b
19. Checking Persistence in Spark UI
The Spark UI is one of the best tools for understanding persistence.
The Storage tab can show information about persisted datasets.
You can inspect things such as:
- Dataset/RDD
- Storage level
- Number of partitions
- Memory usage
- Disk usage
- Fraction of partitions cached
Conceptually:
StorageRDD / Dataset-------------------------------------Storage Level: MEMORY_AND_DISKCached Partitions: 120 / 120Memory Used: 4.2 GBDisk Used: 0 GB-------------------------------------
This is much more useful than blindly calling persist() and assuming it worked.
20. A Real Production Example
Imagine an insurance company processing millions of claims.
We start with:
claims = spark.read.parquet("claims/")
Then we perform expensive processing:
processed_claims = ( claims .filter("claim_amount > 10000") .join(customer_data, "customer_id") .join(policy_data, "policy_id"))
Now we need multiple analytical outputs:
processed_claims.groupBy("state").sum("claim_amount")processed_claims.groupBy("policy_type").sum("claim_amount")processed_claims.groupBy("customer_segment").count()processed_claims.groupBy("risk_category").count()
Without persistence, Spark may repeatedly execute expensive upstream work.
Instead:
processed_claims.persist(StorageLevel.MEMORY_AND_DISK)processed_claims.count()
Then execute the downstream operations.
Finally:
processed_claims.unpersist()
The pattern is:
Read ↓Transform ↓Expensive Join ↓Persist ↓ ├── Analysis 1 ├── Analysis 2 ├── Analysis 3 └── Analysis 4 ↓Unpersist
This is a classic persistence use case.
21. Persistence vs Caching
People often use these terms interchangeably, and in Spark that’s understandable.
The practical relationship is:
Caching ↓Using Spark's default persistence behaviorPersistence ↓Explicitly choosing how the dataset should be stored
Example:
df.cache()
versus:
df.persist(StorageLevel.MEMORY_AND_DISK)
If you need control over the storage strategy, use persist().
22. Persistence vs Checkpointing
This distinction is frequently asked in interviews.
Consider:
Dataset
|
--------------------
| |
Persistence Checkpointing
| |
Reuse computation Break lineage
| |
Performance Fault recovery
Persistence
Used primarily when:
“I will use this dataset repeatedly.”
Checkpointing
Used primarily when:
“I want to materialize this dataset and truncate the lineage.”
Therefore:
persist() ≠ checkpoint()
They solve different problems.
23. A Common Mistake: Calling persist() Everywhere
Consider:
df1.persist()df2.persist()df3.persist()df4.persist()
This looks like optimization.
It isn’t necessarily.
The correct question is:
How many times will this dataset be reused, and how expensive is recomputing it?
A useful mental model is:
Persistence Benefit ≈Cost of recomputation avoided -Cost of storing/managing data
If the recomputation cost is low and storage cost is high, don’t persist.
If the recomputation cost is high and reuse is frequent, persistence becomes attractive.
24. Practical Decision Framework
Before using persist(), ask these five questions:
Question 1
Is this dataset expensive to compute?
If no → probably don’t persist.
Question 2
Will it be used more than once?
If no → probably don’t persist.
Question 3
How large is the dataset?
If huge → carefully select the storage level.
Question 4
Do I have sufficient executor resources?
If no → persistence may hurt performance.
Question 5
When will I stop needing it?
If you know the point, remember to:
df.unpersist()
25. Best Practices
✅ Persist expensive intermediate datasets
Especially when reused by multiple downstream operations.
✅ Choose the storage level deliberately
Don’t automatically assume MEMORY_ONLY is always best.
✅ Materialize intentionally
Use an action after persistence when you want to populate the persisted data before subsequent operations.
✅ Monitor Spark UI
Check whether the dataset actually fits in the selected storage level.
✅ Unpersist when finished
Avoid holding unnecessary data in executor storage.
✅ Avoid persisting everything
Persistence is a performance optimization, not a default coding style.
✅ Benchmark
Compare:
Without persistence vsWith persistence
Measure actual job runtime and resource usage.
26. Complete Example
Here’s a simple production-style example:
from pyspark.sql import SparkSessionfrom pyspark import StorageLevelspark = SparkSession.builder \ .appName("PersistenceExample") \ .getOrCreate()sales = spark.read.parquet("sales/")processed_sales = ( sales .filter("amount > 1000") .join(customers, "customer_id"))# Persist because the dataset is reused multiple timesprocessed_sales.persist(StorageLevel.MEMORY_AND_DISK)# Materialize persisted partitionsprocessed_sales.count()# Reuseregional_sales = ( processed_sales .groupBy("region") .sum("amount"))product_sales = ( processed_sales .groupBy("product") .sum("amount"))regional_sales.show()product_sales.show()# Release persisted dataprocessed_sales.unpersist()
The important part isn’t the syntax.
It’s the reasoning:
Expensive computation ↓Repeated reuse ↓Persist ↓Multiple operations ↓Unpersist
27. Key Takeaways
Persistence is one of Spark’s most important performance optimization techniques.
Remember these points:
persist()stores computed partitions for reuse.cache()uses Spark’s default persistence behavior.persist()allows you to choose a storage level.- Persistence happens at the partition level.
persist()itself doesn’t eagerly compute the dataset.- An action is needed to materialize persisted data.
MEMORY_ONLYcan require recomputation if partitions don’t fit.MEMORY_AND_DISKcan spill partitions to disk.- Persistence consumes executor resources.
- Don’t persist datasets that are cheap or used only once.
- Use Spark UI → Storage to inspect persisted data.
- Call
unpersist()when the dataset is no longer required. - Persistence and checkpointing solve different problems.
- The best storage level depends on workload, data size, and available resources.
The most important rule is simple:
Persist when the cost of recomputation is higher than the cost of storing the intermediate result.
That’s the mindset that separates blindly using Spark optimizations from actually optimizing Spark workloads.
🔗 Continue the Apache Spark Learning Path
← Previous: [Caching in Spark: When It Saves You and When It Burns Your Memory]
Next →: [Checkpointing in Apache Spark: Breaking Lineage and Improving Fault Recovery]
📚 Complete: The Complete Apache Spark Learning Path: Beginner to Advanced
📖 Official Documentation
For the implementation details and supported storage levels, refer to the official Apache Spark documentation:
- Apache Spark RDD Programming Guide — RDD Persistence
- Apache Spark PySpark API — RDD.persist()
- Apache Spark SQL Performance Tuning
If you’re learning Spark seriously, don’t just memorize MEMORY_ONLY and MEMORY_AND_DISK.
Understand why you would choose one over another.
That decision is what matters in a real production Spark application—and it is exactly the kind of reasoning expected in Spark and Data Engineering interviews.