Spark Executor OOM in Shuffle — Memory Math First
Don't double executor memory blind: shuffle OOMs come from skew, fat partitions, or unified-memory pressure.
20+ years shipping production backend systems. Drawn from code that ran under real load.
- ✓Spark memory basics: driver vs executor heaps
- ✓Reading shuffle bytes in the Spark UI task table
- ✓Joins and hash partitioning fundamentals
- Executor OOM during shuffle means execution memory lost its fight inside Spark's unified region — confirm with heap dumps and GC logs before buying RAM
- Unified memory splits heap (minus 300MB reserved) by spark.memory.fraction (0.6); storage can borrow from execution but eviction storms still kill tasks
- Compare the dead task's shuffle bytes to the stage median: a 10x outlier is skew (salt the key), uniform giants mean too few partitions (raise shuffle partitions)
- Off-heap (spark.memory.offHeap.size) and Kryo shrink pressure, but the durable fix is partition math — 128-256MB per partition — plus unpersisting caches you no longer need
Think of each executor as a kitchen with one shared counter. Cooking (execution: sorts, shuffles) and plating (storage: cached data) share it, and cooking may shove plates aside when pressed. A shuffle OOM means the cooks ran out of counter: one drew a 500-plate order while the rest made sandwiches (skew), the kitchen accepted more orders than fit (too few partitions), or old plates never cleared (stale cache). A bigger kitchen helps only case two — split big orders, clear old plates.
Shuffle is where Spark jobs go to die. Maps finish green, the progress bar stalls at 68%, and then executors start dropping with java.lang.OutOfMemoryError one after another. The instinct — double spark.executor.memory and re-run — works just often enough to become habit, and fails just often enough to become expensive. A fleet-wide 8G-to-16G bump on 200 executors doubles your bill on every future run, and when the killer was one skewed key, the same task dies on the bigger box anyway.
Executor memory in Spark isn't one bucket but a negotiated region. The JVM heap minus a 300MB reserve is split by spark.memory.fraction (default 0.6) into unified memory shared between execution (shuffles, joins, sorts) and storage (cached blocks), with execution able to evict storage when pressured. An OOM during shuffle means that negotiation broke down: too much data arrived at one task, the cache refused to yield fast enough, or off-heap relief was never enabled.
This article gives you the memory math first and the sizing second. You'll learn to read unified-memory boundaries, to separate skew deaths from genuine capacity deaths with one histogram, and to apply the real fixes — salting, partition math, off-heap, Kryo, cache hygiene — in the order that actually holds. Size the box last, not first.
Unified Memory: The 60% Region That Decides Everything
Spark's executor heap is divided before your data ever arrives. First, 300MB is reserved off the top (spark.memory.reserved — untouchable, pays for internal objects). Of what remains, spark.memory.fraction (default 0.6) becomes the unified region shared by execution and storage; the other 40% belongs to user code, UDFs, and internal metadata. Inside the unified 60%, execution memory (shuffles, sorts, joins) and storage memory (cached blocks, controlled by spark.memory.storageFraction, default 0.5 of the unified region) borrow from each other — with one iron rule: execution can evict storage, but storage can never evict execution.
That asymmetry is the whole ballgame for shuffle OOMs. When a giant reduce task arrives, execution demands space and cached blocks get evicted — fine, unless the cache is referenced by running tasks (then it can't evict cleanly) or eviction storms coincide with peak sort pressure. The task needs its sort workspace now, the cache drains slowly, and the JVM hits the ceiling mid-negotiation. Your stderr says OutOfMemoryError; the real story is a borrowing protocol that couldn't keep up with a lopsided demand spike.
Run the numbers for a concrete 8G executor: 8192MB minus 300MB reserve leaves 7892MB; 60% unified gives ~4735MB; storageFraction 0.5 initially reserves ~2367MB for cached blocks. A single 60GB hot partition (like our incident's 1TB key at finer partitioning) doesn't fit in 4735MB no matter how you tune the fractions — which is why fraction-tuning only helps borderline cases and never fixes skew. Measure the region, respect its ceiling, and fix demand (partition size) before supply (heap).
Skew vs Genuine Size: The One Ratio That Decides
Every shuffle OOM is either skew (a few tasks drowning while siblings idle) or genuine size (all tasks uniformly too fat) — and the response to each is the opposite of the other. Skew demands parallelism surgery: salt keys, split partitions, broadcast the small side. Genuine size demands arithmetic: more partitions, leaner rows, bigger boxes. Apply the skew fix to a sizing problem and you add complexity for nothing; apply the sizing fix to skew and you burn money growing boxes around one immortal task.
The deciding measurement is the top-to-median ratio of shuffle-read bytes in the dead stage's task table. Sort descending, note the largest task, note the median, divide. Above ~10x is skew with near certainty — no uniform distribution produces that gap, and no memory setting bridges it. Near 1x with every task above ~1GB is genuine size: your spark.sql.shuffle.partitions default of 200 is carving terabytes into door-sized slabs and each executor is genuinely being asked to swallow too much.
Confirm skew's source with a key census before choosing the weapon. GROUP BY the join key and check the top share: a single key over ~10% of rows needs salting (or AQE skew splitting), while a heavy-hitter band of thousands of moderately hot keys responds better to broadcast hints or two-stage aggregation. The ratio triages in a minute; the census prescribes in five.
Salting Hot Keys: Splitting the Unsplittable Task
Salting breaks one hot key into N sub-keys that hash to N different partitions, turning an atomic monster task into N ordinary ones. On the large side you append a random suffix (0 to N-1); on the small side you replicate each row N times, once per suffix; you join on the salted key, then strip the suffix and re-aggregate. With N=32, our incident's 1TB hot partition became 32 tasks of ~32GB shuffle each — still chunky, but each fits comfortably in an 8G executor's streaming sort workspace since sorts spill to disk incrementally.
Pick N from the hot key's bytes: N = hot_key_bytes / 200MB, rounded up to a power-friendly number (16, 32, 64). Too small and sub-tasks still OOM; too large and you pay replication overhead on the small side (N copies of every small-side row). When only one side is skewed, salt just that side's key and broadcast or replicate the other — never salt both sides symmetrically unless both census hot.
Spark 3 offers two automated alternatives worth trying before hand-salting. The SKEW join hint (SELECT /+ SKEW('events', 'device_id') / ...) tells the optimizer which key skews so it can split partitions for you, and Adaptive Query Execution's skew handling (spark.sql.adaptive.skewJoin.enabled, on by default since Spark 3.2 alongside AQE) detects skewed shuffle partitions at runtime and splits them automatically. Both fail on extreme skew (a single key worth 11% of the data can defeat the defaults) — tune spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes (256MB) and skewedPartitionFactor (5) before falling back to manual salt.
Partition Counts: Landing Every Slice at 128-256MB
spark.sql.shuffle.partitions defaults to 200 — a value chosen when terabytes were exotic. At 9.4TB of shuffle, 200 partitions means 47GB per task, and no executor heap on earth sorts 47GB in memory. The rule that replaces the default: total shuffle bytes divided by ~200MB, rounded up, is your partition count. Our 9.4TB needed roughly 49,000 partitions; the job ran 2,000. That single misconfiguration was responsible for every 'uniform giant' OOM in the fleet's history.
Set it per-query rather than globally when stages differ wildly: spark.conf.set('spark.sql.shuffle.partitions', '8192') before the big join, smaller values for small stages (too many partitions on small data means task-scheduling overhead dominates — thousands of 2-second tasks with 1-second launch tax). For RDD code the equivalents are the numPartitions argument to repartition/join and spark.default.parallelism (defaults to total cores, a better starting point than 200 but still worth explicit math).
Verify, don't trust: after one test run, check the task table's shuffle-read column — p99 should sit 128-256MB. Smaller wastes scheduling overhead and explodes the driver's task bookkeeping; larger risks heap-busting giants. AQE's coalescing (spark.sql.adaptive.coalescePartitions.enabled, on by default) can auto-merge small post-shuffle partitions at runtime, which protects the small side while your explicit count protects the big side — belt and suspenders that actually compose.
Off-Heap, Kryo, and Cache Hygiene: Shrinking Pressure
When partitions are healthy and code is lean but tasks still die near the ceiling, shrink Spark's per-byte overhead instead of growing the box. Three levers stack: serialization format, off-heap sort workspace, and cache lifecycle. Together they've saved jobs that partition math alone couldn't — typically the ones with wide rows, cached inputs feeding multiple actions, and Python UDFs churning objects.
Kryo serialization (spark.serializer=org.apache.spark.serializer.KryoSerializer) serializes shuffle data 2-10x more compactly than Java serialization with far less CPU, directly shrinking both shuffle bytes and deserialization garbage. Register custom classes with spark.kryo.registrationRequired=true in dev to catch unregistered fallbacks. Off-heap execution memory (spark.memory.offHeap.enabled=true plus spark.memory.offHeap.size, e.g. 6G) moves Tungsten sort and aggregation buffers outside the GC-managed heap — sorts that used to trigger full-GC death spirals now spill to off-heap and disk without pausing the world.
Cache hygiene is the silent killer because it's invisible in the task table. Every persist() holds unified memory until unpersist(), and a DataFrame cached for stage 2 still occupies storage memory during stage 9's shuffle — memory execution could have borrowed but can't reclaim fast enough under spike pressure. Audit with the Storage tab: any RDD/DataFrame with 100% cached fraction that no running action needs is stolen execution memory. Unpersist the moment the last dependent action completes, prefer MEMORY_AND_DISK over MEMORY_ONLY (spill beats eviction storms), and never cache the input of a single-pass ETL.
persist() as a loan you repay with unpersist() right after the last dependent action — the Storage tab showing stale 100%-cached frames during a shuffle OOM is the smoking gun.Sizing the Box Last: Heap, Overhead, and Cores
Only after skew is salted, partitions land 128-256MB, Kryo is on, caches are freed, and tasks still die flat against the ceiling is the box genuinely too small. Then size deliberately: raise spark.executor.memory one step (8G→12G, not 8G→32G) and watch whether the failure mode changes — a moved corpse (dies later, same signature) means remaining skew, while a cured stage means true capacity. Step sizing costs one test run; leap sizing costs a permanent 4x bill.
Read YARN's kill message as two separate diagnoses. 'Container killed for exceeding physical memory' with heap healthy means overhead was thin: native code, off-heap buffers, and Python workers all live outside heap, and the default max(384M, 10% of heap) overhead doesn't cover shuffle-heavy Python stages — set spark.executor.memoryOverhead to 2-4G explicitly. 'Java heap space' with overhead healthy means the JVM truly filled: add heap, reduce cores per executor (fewer concurrent tasks share the heap — spark.executor.cores 5→2 halves per-task contention), or both.
Cores are the forgotten memory knob. Five concurrent tasks on one 8G executor share ~4.7G of unified memory (~950MB each); two concurrent tasks get ~2.4G each. Dropping cores per executor is often cheaper than raising memory because cloud pricing charges RAM whether tasks share it well or not. Final guardrail: cap any sizing change with a cost-per-run alert so emergency settings can't silently become permanent — our incident's 32G would have taxed every nightly run 4x forever if nobody reverted it.
The $12,000 Shuffle: Doubling Memory Twice for One Hot Key
rand()*32).cast('int')) on the big side, exploded with the 32 suffixes on the small side, then second-stage aggregation stripped the suffix. Shuffle partitions went 2,000→4,096 to land partitions near 200MB, Kryo replaced Java serialization, and executor memory returned to 8G. Runtime dropped from 7 dying hours to 94 green minutes, and the bill fell 61% versus the 32G attempt — salting cost one day of work against three weeks of doublings.- One task 50x slower than the median is never a sizing problem — it's an atomicity problem no box can fix. Read the task-duration histogram before the pricing page; the 5-hour outlier convicted skew on day one.
- Memory scales the box, salting scales the parallelism. Doubling RAM leaves one task holding 1 TB; 32 salt buckets turn it into 32 tasks holding 32 GB each — the only fix that changes the shape of the work.
- Revert emergency sizing after the real fix. The 32G setting would have taxed every future run forever; returning to 8G after salting locked in a 61% cost cut instead of a permanent surcharge.
| File | Command / Code | Purpose |
|---|---|---|
| unified_memory_math.py | heap_mb = 8 * 1024 | Unified Memory |
| skew_vs_size_census.sql | SELECT device_id, COUNT(*) AS n, | Skew vs Genuine Size |
| salted_skew_join.py | from pyspark.sql import functions as F | Salting Hot Keys |
| PartitionMath.scala | val shuffleBytes = 9.4 * 1024L * 1024L * 1024L * 1024L // 9.4 TB | Partition Counts |
| lean_shuffle_conf.py | spark.conf.set("spark.serializer", | Off-Heap, Kryo, and Cache Hygiene |
Key takeaways
persist() is a loanCommon mistakes to avoid
5 patternsDoubling executor memory fleet-wide for a skewed stage
Tuning spark.memory.fraction and storageFraction blindly
Leaving spark.sql.shuffle.partitions at the 200 default on terabyte shuffles
Caching inputs and forgetting them across a multi-stage job
Serializing shuffles with default Java serialization plus Python UDFs
Interview Questions on This Topic
An executor dies with OutOfMemoryError mid-shuffle. What's your first measurement?
Frequently Asked Questions
20+ years shipping production backend systems. Drawn from code that ran under real load.
That's Spark. Mark it forged?
6 min read · try the examples if you haven't