Sathus AI 2.0 is now generally available — evaluation harnesses and guardrails included. Explore
The definitive engineering triage guide to diagnosing and fixing Apache Spark OutOfMemory (OOM) crashes. Demystifying driver heap exhaustion from collect() vs executor container kills (Exit code 137), broadcast hash join thresholds, shuffle skew, Adaptive Query Execution (AQE), and off-heap memory overhead.
PySpark OutOfMemory (OOM) errors diverge fundamentally by failure boundary: Driver OOM (java.lang.OutOfMemoryError: Java heap space) occurs when client-side operations like .collect(), .toPandas(), broadcast joins exceeding spark.sql.autoBroadcastJoinThreshold, or bloated DAG accumulators pull distributed data onto a single JVM. Conversely, Executor OOM (Exit code 137 / SIGKILL or "Container killed by YARN/K8s for exceeding memory limits") occurs when executor execution memory cannot accommodate high shuffle partition skew, unvectorized Python UDF worker processes, or off-heap overhead. Remediation requires isolating the Spark UI Stage failure point, replacing .collect() with windowed partitions, enforcing 100MB-200MB target partition sizing via AQE, salting skewed join keys, and provisioning spark.executor.memoryOverhead to at least 15-20% of executor heap.
The definitive engineering triage guide to diagnosing and fixing Apache Spark OutOfMemory (OOM) crashes. Demystifying driver heap exhaustion from collect() vs executor container kills (Exit code 137), broadcast hash join thresholds, shuffle skew, Adaptive Query Execution (AQE), and off-heap memory overhead.
A: Snappy-compressed Parquet files are heavily columnar compressed. When loaded into JVM memory and converted into a broadcast hash table with Java object pointer overhead, 500MB of compressed storage can expand to 3.5GB to 5GB of raw in-memory Java objects, blowing past the default 10MB broadcast threshold and exhausting executor heap.
A: We recommend 4 to 5 cores per executor with 16GB to 28GB of memory. Allocating more than 5 cores leads to severe JVM garbage collection pauses, while fewer than 3 cores wastes executor container overhead.
A: spark.memory.fraction (default 0.60) defines the overall pool of JVM heap allocated to Spark (the remaining 40% is user memory for custom objects). spark.memory.storageFraction (default 0.50 of the fraction) defines the immune storage threshold: execution can borrow from storage, but cached RDDs cannot evict active execution memory.
Understanding Spark memory requires separating the Driver JVM from Executor container boundaries: 1. Driver Memory: Responsible for DAG scheduling, SparkContext, Catalyst query optimization, and holding results passed back to the Python client. When code invokes .collect() or .toPandas(), all distributed partitions are transferred over the network into the Driver JVM. If the serialized data exceeds spark.driver.memory, the driver throws java.lang.OutOfMemoryError: Java heap space and terminates the entire application. 2. Executor Container Architecture: An executor runs inside a YARN container or Kubernetes pod governed by two memory pools: • On-Heap JVM Memory (spark.executor.memory): Subdivided into Reserved Memory (300 MB), User Memory (40%), and Unified Memory (60% divided dynamically between Execution and Storage). • Off-Heap / Overhead Memory (spark.executor.memoryOverhead): Crucial in PySpark because Python worker processes (PyArrow, Pandas UDFs, NumPy) execute outside the JVM. If Python processes + JVM overhead exceed the container limit, the Linux kernel OOM killer or Kubernetes cgroups issues SIGKILL (Exit code 137).
Two primary culprits cause driver OOM crashes: • Unbounded Collection: Calling df.collect() or df.toPandas() on a 50GB dataset when driver memory is 8GB. Even df.take(1000) can crash the driver if individual rows contain massive nested JSON strings or multi-megabyte embeddings. • Broadcast Join Memory Expansion: Spark broadcast joins (spark.sql.autoBroadcastJoinThreshold, default 10MB) seem safe on disk. However, a 10MB Snappy-compressed Parquet file can expand to 150MB-300MB of uncompressed Java objects inside the driver memory before being serialized and broadcast to executors. If 3 broadcasts run concurrently, the driver crashes.
Executor crashes usually happen during Shuffle Hash Joins or Sort-Merge Joins: • The Null Key Skew Trap: If an e-commerce dataset has 100 million records where 30% of user_id values are NULL or "GUEST", the default HashPartitioner maps all those rows to a single partition. While 199 tasks finish in seconds, task #200 processes 30GB alone on one executor, exceeding execution memory and crashing. • Spill (Memory) vs Spill (Disk): When execution memory fills up, Spark spills intermediate shuffle data to local disk. If disk space is exhausted, or the serialized buffer exceeds 2GB (the Spark ByteBuffer integer limit), the executor throws org.apache.spark.shuffle.FetchFailedException.
Apply these enterprise configuration formulas to permanently stabilize pipelines: 1. Partition Sizing Formula: Target partition count = max(200, Stage Input Data Size in MB / 128 MB). Keep partitions between 100MB and 200MB. 2. Overhead Memory Rule for PySpark: Set spark.executor.memoryOverhead to max(384m, 0.20 * spark.executor.memory). For heavy PyArrow/Pandas UDF workloads, increase to 0.25. 3. Adaptive Query Execution (AQE): Ensure spark.sql.adaptive.enabled=true, spark.sql.adaptive.skewJoin.enabled=true, and spark.sql.adaptive.coalescePartitions.enabled=true.
Distributed Spark memory boundary separating Driver JVM, Executor JVM on-heap pool, and off-heap OS container limits.
Coordinates DAG, Catalyst optimization, and broadcast tables. Sized via spark.driver.memory.
Dynamic 60/40 pool sharing between Execution (shuffles/joins) and Storage (caches).
Dedicated memory for PySpark Python workers, PyArrow C++ buffers, and NIO shuffle buffers.
Pinpoint root cause failure modes and match observed metrics to actionable remediation.
| UI Tab / Tool | Observed Metric / Signal | Underlying Failure Mode | Actionable Remediation |
|---|---|---|---|
| Spark UI: Executors Tab | Executors marked "Dead" with Exit Code 137 / SIGKILL | Off-heap memory breach. Python worker processes (PyArrow/Pandas) or JVM GC overhead exceeded container memory allocation. | Increase spark.executor.memoryOverhead to 20-25% of spark.executor.memory; reduce cores per executor from 8 to 4. |
| Spark UI: Stages Tab | Task duration distribution shows Max time 45m while 75th percentile is 12s | Extreme data skew. A single join or group-by key concentrated enormous volume onto one executor task. | Enable spark.sql.adaptive.skewJoin.enabled=true or implement two-phase key salting on the join key. |
| Spark UI: Stage Task Metrics | Spill (Memory) > 50 GB and Spill (Disk) > 15 GB on shuffle stage | Execution memory pool saturated during SortMergeJoin, triggering aggressive disk serialization. | Double spark.sql.shuffle.partitions and tune spark.memory.fraction to 0.70. |
| Driver Log / stderr | java.lang.OutOfMemoryError: Java heap space during collect() | Client action attempted to serialize distributed DataFrame into local driver JVM. | Replace .collect() with .take(100), write directly to Delta/S3, or paginate using row_number() windows. |
import pyspark.sql.functions as F
from pyspark.sql import SparkSession
def join_skewed_dataframes(large_skewed_df, dimension_df, join_key, num_salts=16):
"""
Eliminates executor OOM by salting skewed join keys across cluster:
1. Adds random salt suffix (0..num_salts-1) to skewed large table
2. Replicates dimension table keys across all salt values using explode
3. Executes uniform hash join without single-node hotspotting
"""
# 1. Salt the skewed large DataFrame
salted_large_df = large_skewed_df.withColumn(
"salt_val",
F.floor(F.rand() * num_salts)
).withColumn(
"salted_join_key",
F.concat(F.col(join_key), F.lit("_"), F.col("salt_val"))
)
# 2. Replicate dimension DataFrame with an array of all possible salt values
salt_array = F.array([F.lit(i) for i in range(num_salts)])
replicated_dim_df = dimension_df.withColumn("salt_val", F.explode(salt_array)).withColumn(
"salted_join_key",
F.concat(F.col(join_key), F.lit("_"), F.col("salt_val"))
)
# 3. Perform balanced join on salted key
joined_df = salted_large_df.join(
replicated_dim_df,
on="salted_join_key",
how="inner"
).drop("salt_val", "salted_join_key")
return joined_dfA 16-node cluster running default settings continually failed at Stage 4 with Exit Code 137 due to null-key skew. Manual restarts and ad-hoc partition scaling required 5h 15m runtime, spilling 140 GB to disk, and crashed the driver on broadcast joins.
Configured 128MB partition sizing, AQE skew join splitting, 20% off-heap memoryOverhead allocation, and a salted two-phase join pattern. Pipeline runtime slashed from 5h 15m to 32 minutes with zero disk spills and zero container kills under synthetic benchmark load.
Simulated benchmark scenario comparing unoptimized Spark 3.5 default partition sizing against AQE + two-phase salting on a modeled 2TB skewed dataset with 30% key null-concentration. Runtime, container survival, and disk spill metrics reflect engineering simulation parameters.
Compute container allocation, off-heap buffers, and AQE partition targets based on Sathus distributed systems sizing formulas.
# Generated by Sathus Distributed Systems Calculator
spark.executor.instances: 16
spark.executor.cores: 4
spark.executor.memory: 10g
spark.executor.memoryOverhead: 3584m
spark.sql.shuffle.partitions: 2000
spark.sql.adaptive.enabled: true
spark.sql.adaptive.skewJoin.enabled: true
spark.sql.adaptive.coalescePartitions.enabled: true
spark.memory.fraction: 0.60
spark.memory.storageFraction: 0.50| Failure Mode | Error Signature | Root Cause | Immediate Triage Fix | Permanent Architectural Fix |
|---|---|---|---|---|
| Driver Heap OOM | java.lang.OutOfMemoryError: Java heap space | .collect(), .toPandas(), broadcast too large | Increase spark.driver.memory to 16g | Write directly to storage, paginate reads |
| Executor SIGKILL | Exit code 137 / Container killed by YARN | Off-heap PyArrow/Python worker breach | Increase spark.executor.memoryOverhead | Convert Python UDFs to PySpark native/Pandas UDFs |
| Data Skew OOM | Stage hangs at 99%, single task fails OOM | Non-uniform key distribution (nulls/hot keys) | Filter null keys before join | Implement two-phase key salting + enable AQE |
| Disk Spill Exhaustion | FetchFailedException: Connection reset by peer | Shuffle data exceeds execution pool and disk | Increase spark.sql.shuffle.partitions | Tune partition sizing to 128MB per task |
Snappy-compressed Parquet files are heavily columnar compressed. When loaded into JVM memory and converted into a broadcast hash table with Java object pointer overhead, 500MB of compressed storage can expand to 3.5GB to 5GB of raw in-memory Java objects, blowing past the default 10MB broadcast threshold and exhausting executor heap.
We recommend 4 to 5 cores per executor with 16GB to 28GB of memory. Allocating more than 5 cores leads to severe JVM garbage collection pauses, while fewer than 3 cores wastes executor container overhead.
spark.memory.fraction (default 0.60) defines the overall pool of JVM heap allocated to Spark (the remaining 40% is user memory for custom objects). spark.memory.storageFraction (default 0.50 of the fraction) defines the immune storage threshold: execution can borrow from storage, but cached RDDs cannot evict active execution memory.
Principal Distributed Systems Engineer
Part of the Big Data & Spark Engineering at Sathus Technology. Specializing in mission-critical data lakehouses, streaming analytics, and compliance-driven platforms.
Book a Spark Cluster & Pipeline Audit with Sathus Distributed Systems Architects to eliminate OOM failures and reduce cloud compute costs by up to 60%.