Spark持久化机制诡异问题求助:大Dataset持久化异常
Hey there, let’s tackle this weird persistence problem you’re hitting with your large Dataset. I’ve dealt with similar headaches on older Spark versions like 2.1.1, so let’s break down what might be going on and how to fix it.
First, Let’s Recap Your Setup
Just to make sure I’m on the same page:
- Spark Version: 2.1.1 (important—this release has some known quirks with persistence)
- Cluster Specs: 9 nodes, 100GB RAM/15 cores each (total cluster RAM: 900GB, total cores: 135)
- App Config: 27 executors (3 per node), 30GB RAM/5 cores each (total executor RAM: 810GB, total cores: 135—you’re fully utilizing the cluster)
- Dataset: 810 partitions, in-memory size ~1.4TB (way larger than total executor RAM, so
MEMORY_AND_DISK_SERmakes sense here to spill excess to disk)
Likely Causes of the "Weird" Behavior
Since you didn’t share the exact error message, I’ll cover the most common culprits for this scenario:
1. Spark 2.1.1’s MEMORY_AND_DISK_SER Bugs
Older Spark versions (pre-2.3) had several known issues with this persistence level:
- Incorrect calculation of serialized data size, leading to unexpected OOMs even when disk should be used
- Occasional data corruption or loss when spilling to disk
- Inefficient spill logic that causes excessive I/O and slowdowns
2. Uneven Partition Sizes
Your 810 partitions sound evenly calculated, but if some partitions are way larger than the average (~1.7GB/partition), you could run into problems. Spark can’t split a single partition across memory and disk—if a serialized partition exceeds an executor’s available storage memory, it’ll throw an OOM.
3. Insufficient Disk Space
You’re looking to spill roughly 590GB (1.4TB - 810GB) to disk. With 3 executors per node, each node needs ~65GB of free space. If any node is low on disk, the spill will fail with a "No space left on device" error.
4. Inefficient Serialization (Java Default)
Spark’s default Java serializer is slow and produces larger serialized data. If your Dataset has complex objects, the serialized size could be 2-3x larger than the in-memory size, pushing you over memory/disk limits faster than expected.
5. Memory Competition Between Storage and Execution
Spark 2.x uses dynamic memory allocation—storage memory (for persistence) and execution memory (for tasks) share a pool. With 5 cores per executor, you’re running 5 concurrent tasks per executor, which might be eating into the storage memory pool, forcing more spills than necessary or causing OOMs.
Step-by-Step Troubleshooting & Fixes
1. First: Get the Exact Error Logs
This is critical. Check the driver and executor logs for specific errors (e.g., OutOfMemoryError, disk space errors, or serialization exceptions). The error message will narrow down the root cause instantly.
2. Check Partition Size Uniformity
Use these quick checks to see if you have oversized partitions:
- Run this code to get partition sizes:
import org.apache.spark.sql.functions.spark_partition_id df.groupBy(spark_partition_id()).count().orderBy("count").show(false) - Or look at the Storage tab in the Spark UI—you can see the size of each persisted partition.
If you find large partitions, re-partition the Dataset to a higher number (e.g., 1620 partitions) to make each one smaller:
val repartitionedDf = df.repartition(1620) repartitionedDf.persist(org.apache.spark.storage.StorageLevel.MEMORY_AND_DISK_SER)
3. Switch to Kryo Serialization
Kryo is faster and produces much smaller serialized data. Configure it in your app:
spark.conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer") spark.conf.set("spark.kryo.registrationRequired", "false") // Skip class registration for quick testing; register classes later for better performance
This can reduce serialized size by 50% or more, easing memory and disk pressure.
4. Adjust Memory Allocation Settings
- Increase Storage Memory Ratio: If your app doesn’t need heavy execution memory, tweak the storage fraction:
This increases the portion of the memory pool reserved for storage (from 50% to 75% of the shared pool), giving more room for persisted data.# Add this to your spark-submit command --conf spark.memory.storageFraction=0.75 - Reduce Executor Cores: Try lowering executor cores to 4 instead of 5. This gives each core more RAM (30GB/4 = 7.5GB vs 6GB) and reduces concurrent task memory competition:
--num-executors 33 --executor-cores 4 --executor-memory 27GB # Keeps total cores/RAM roughly the same
5. Verify Disk Space
Check free disk space on all cluster nodes. If space is tight:
- Clean up old logs or unused files
- Mount additional disk storage for Spark’s local directories (configure with
spark.local.dir)
6. Consider Upgrading Spark (Long-Term Fix)
Spark 2.1.1 is over 6 years old. Versions 2.3+ fixed many persistence-related bugs, improved spill logic, and added better memory management. If your infrastructure allows it, upgrading will prevent many of these issues from recurring.
内容的提问来源于stack exchange,提问作者Malkolm

