Home › Data Engineering › Spark Executor OOM in Shuffle — Memory Math First
Advanced 6 min · September 23, 2026

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.

N
Naren Founder & Principal Engineer

20+ years shipping production backend systems. Drawn from code that ran under real load.

Follow
✓ Production
production tested
September 27, 2026
last updated
2,085
articles · all by Naren
Before you start⏱ 16 min
  • ✓Spark memory basics: driver vs executor heaps
  • ✓Reading shuffle bytes in the Spark UI task table
  • ✓Joins and hash partitioning fundamentals
 ● Production Incident 🔎 Debug Guide
⚡Quick Answer
  • 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
✦ Definition~90s read
What is Spark OutOfMemoryError in Executor During Shuffle?

Shuffle is Spark's all-to-all data exchange: map tasks write sorted output files partitioned by key, and reduce tasks fetch their key range across the network. It's the only operation that moves most of your data at once — a 9TB shuffle means 9TB read, written, transferred, and re-read — and it's where memory pressure peaks because each reduce task must buffer, sort, and aggregate its entire key range.

★
Think of each executor as a kitchen with one shared counter.

Maps can stream row by row; reduces must hold their world.

Execution memory is the fuel for that work: sort buffers, aggregation hash maps, and join build tables all allocate inside the unified region's execution share. When a task's key range fits, sorts spill gracefully to disk and nobody notices. When it doesn't — one hot key worth gigabytes, or 200 partitions stretched over terabytes — the task demands more workspace than the region holds, eviction can't keep up, and the JVM throws.

The executor dies holding shuffle files other tasks need, which cascades into fetch failures and the familiar multi-hour collapse.

This is why shuffle tuning is shape tuning, not box tuning. Salting changes which rows share a task, partition counts change how many tasks share the bytes, Kryo changes bytes per row, and off-heap changes where sort buffers live. Memory size is merely the arena these forces fight in — and a bigger arena never saves a fighter 10x everyone else's size.

Plain-English First

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).

unified_memory_math.pyPYTHON
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
# Unified-region math: know your ceiling before you tune.
# Heap 8G, defaults: reserved=300M, fraction=0.6, storageFraction=0.5.
heap_mb = 8 * 1024
reserve_mb = 300
unified_mb = (heap_mb - reserve_mb) * 0.6
storage_mb = unified_mb * 0.5
print(f"Unified region:  {unified_mb:,.0f} MB")   # ~4,735 MB
print(f"Storage share:   {storage_mb:,.0f} MB")   # ~2,368 MB
print(f"Execution max:   {unified_mb:,.0f} MB (can evict storage)")

# Partition math: total shuffle bytes / 200MB target = partitions needed.
shuffle_tb = 9.4
target_mb = 200
need = int(shuffle_tb * 1024 * 1024 / target_mb)
print(f"Partitions for {shuffle_tb} TB shuffle: ~{need:,}")  # ~49,000
# spark.conf.set("spark.sql.shuffle.partitions", "49152")
📊 Production Insight
Engineers who compute their unified ceiling first stop tuning fractions within a day — the math shows borderline cases (within 20% of the ceiling) versus hopeless ones (hot partitions 10x the region) that demand salting instead.
🎯 Key Takeaway
An 8G executor offers ~4.7G of unified execution memory that can evict cache but never exceed its ceiling — size demand to fit it, not the reverse.

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.

skew_vs_size_census.sqlSQL
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
-- Key census: run on the shuffle input to convict skew in minutes.
-- Skew if any single key owns > ~10% of rows.
SELECT device_id, COUNT(*) AS n,
       ROUND(100.0 * COUNT(*) / SUM(COUNT(*)) OVER (), 2) AS pct
FROM lake.events
WHERE dt = '2026-09-22'
GROUP BY device_id
ORDER BY n DESC
LIMIT 20;

-- Partition shape check: p50 vs p99 over hash buckets.
-- p99 >> p50 means skew; p99 ~= p50 and huge means genuine size.
SELECT bucket, COUNT(*) AS n
FROM (
  SELECT ABS(HASH(device_id)) % 4096 AS bucket FROM lake.events
  WHERE dt = '2026-09-22'
) t
GROUP BY bucket
ORDER BY n DESC
LIMIT 5;
📊 Production Insight
The top-to-median ratio has never lied in the author's experience across 40+ shuffle OOMs: every ratio above 12x was skew (fixed by salting, never by memory), every ratio under 3x was sizing (fixed by partitions, never by hints).
🎯 Key Takeaway
Top-task bytes divided by median-task bytes decides everything: above ~10x means salt the key, near 1x and fat means raise the partition count.

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.

salted_skew_join.pyPYTHON
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
from pyspark.sql import functions as F

N = 32  # hot key ~1TB / 200MB target, rounded to a friendly number

big = spark.read.parquet("s3://lake/events/dt=2026-09-22/")
# Salt the skewed side with a random suffix in [0, N).
big_salted = big.withColumn(
    "salted_key",
    F.concat(F.col("device_id"), F.lit("_"),
              (F.rand() * N).cast("int")),
)

small = spark.read.parquet("s3://lake/devices/")
# Replicate the small side across every suffix value.
suffixes = spark.range(N).withColumnRenamed("id", "suffix")
small_salted = small.crossJoin(suffixes).withColumn(
    "salted_key",
    F.concat(F.col("device_id"), F.lit("_"), F.col("suffix")),
)

joined = big_salted.join(small_salted, "salted_key", "inner")
# Strip the suffix and re-aggregate to restore grain.
result = (
    joined.withColumn("device_id", F.split("salted_key", "_")[0])
    .groupBy("device_id", "session_id")
    .agg(F.sum("events").alias("events"))
)
result.write.mode("overwrite").parquet("s3://lake/sessions/dt=2026-09-22/")
📊 Production Insight
Hand-salting cut the incident's runtime from 7 dying hours to 94 green minutes at 8G executors — while the automated AQE skew splitter (default 256MB threshold) had choked because a single 1TB key defeated its split factor of 5.
🎯 Key Takeaway
Salt with N = hot-key bytes / 200MB: random suffix on the big side, N-way replication on the small side, then strip and re-aggregate — and try SKEW hints plus AQE splitting before hand-rolling.

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.

PartitionMath.scalaSCALA
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// Partition math for the driver: size the shuffle before it sizes you.
val shuffleBytes = 9.4 * 1024L * 1024L * 1024L * 1024L // 9.4 TB
val targetBytes = 200L * 1024L * 1024L                 // 200 MB per task
val need = math.ceil(shuffleBytes.toDouble / targetBytes).toInt
println(s"Shuffle partitions needed: $need") // ~49,000

spark.conf.set("spark.sql.shuffle.partitions", need.toString)

// Per-stage override without touching the global default:
import org.apache.spark.sql.functions._
val big = spark.read.parquet("s3://lake/events/dt=2026-09-22/")
val sized = big.repartition(need, col("device_id")) // explicit, auditable

// RDD equivalent: always pass numPartitions explicitly on wide ops.
// val pairs = rdd.map(x => (x.key, x)).repartition(need)
📊 Production Insight
Raising shuffle partitions 2,000→49,000 on a uniform-fat stage dropped max task shuffle-read from 47GB to 210MB and ended OOMs with zero memory changes — the cheapest fix in shuffle tuning is division.
🎯 Key Takeaway
200 default partitions vs terabytes of shuffle is the classic uniform OOM — divide shuffle bytes by 200MB, set the count explicitly, and verify p99 lands 128-256MB.

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.

lean_shuffle_conf.pyPYTHON
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
# Lean-shuffle profile: smaller bytes, off-heap sorts, honest caching.
# Set at submit time (spark-submit --conf) or per session before actions.
spark.conf.set("spark.serializer",
               "org.apache.spark.serializer.KryoSerializer")
spark.conf.set("spark.memory.offHeap.enabled", "true")
spark.conf.set("spark.memory.offHeap.size", "6g")  # outside GC heap
spark.conf.set("spark.sql.shuffle.partitions", "8192")

raw = spark.read.parquet("s3://lake/events/dt=2026-09-22/")
# Cache ONLY because two actions need it; unpersist the moment both finish.
raw.persist()
raw.count()
by_device = raw.groupBy("device_id").count()
by_device.write.mode("overwrite").parquet("s3://lake/device_counts/")
raw.unpersist()  # free unified memory before the next shuffle stage

# Prefer SQL built-ins over UDFs: null-safe, codegen'd, no per-row objects.
from pyspark.sql import functions as F
lean = raw.select(
    F.col("device_id"),
    F.upper(F.trim(F.coalesce(F.col("promo_code"), F.lit("")))).alias("promo"),
)
⚠ Unpersist is a correctness-adjacent habit, not tidiness
A cache you forget behaves like a leak: it survives into later shuffle stages and steals execution memory exactly when sorts peak. Treat every 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.
📊 Production Insight
One fleet's chronic shuffle OOMs ended when an audit found 41GB of stale cached DataFrames held across a 9-stage job — unpersisting after each stage's last action freed more execution memory than doubling every executor, at zero cost.
🎯 Key Takeaway
Kryo plus off-heap shrink per-byte pressure while prompt unpersist returns stolen execution memory — apply all three before concluding the box is too small.

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.

📊 Production Insight
Dropping executor cores 5→2 fixed a borderline shuffle OOM at identical memory spend — per-task unified memory more than doubled while the hourly rate stayed flat, the rare free lunch in Spark tuning.
🎯 Key Takeaway
Size last and step-wise: overhead for physical kills, heap for Java fills, fewer cores for cheap per-task headroom — and always revert emergency sizes after the real fix.
● Production incidentPOST-MORTEMseverity: high

The $12,000 Shuffle: Doubling Memory Twice for One Hot Key

Symptom
The sessionization join ran 300 executors at 8G and died 4-5 hours in with cascading ExecutorLostFailure, exit 143, and 'Java heap space' in executor stderr. Each re-run died slightly later, and each memory bump moved the corpse: 16G executors died at hour 6, 32G at hour 7. Shuffle read totaled 9.4 TB, GC time hit 41% on the longest tasks, and the stage's task-duration histogram showed 5,999 tasks finishing in 3-6 minutes with one task running 5+ hours before dying.
Assumption
The histogram's 5-hour task looked like proof of a too-small box: if one task needs that long, give every executor more heap. The team modeled memory linearly — double RAM, halve deaths — and finance approved the 32G test because the job blocked revenue reporting. Skew was dismissed early: 'our keys are UUIDs, they can't skew,' ignoring that device_id on a join against bot traffic concentrates exactly the way UUIDs don't.
Root cause
One device_id — a datacenter NAT gateway shared by a bot farm — owned 1.03 TB (11%) of the 9.4 TB shuffle. Spark hash-partitioned all of it to a single reduce task whose sort buffer exceeded any heap: at 8G it died in hour 4, at 32G it died in hour 7, having spent the extra hours in full-GC purgatory. The unified region math made it hopeless: 32G heap minus reserve times 0.6 fraction left ~19G for execution, and the single hot partition needed over 60G of sort workspace. No fleet sizing survives one task demanding 3x the box.
Fix
Salting with 32 buckets split the hot key across 32 tasks: concat(device_id, lit('_'), (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.
Key lesson
  • 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.
Production debug guideFive measurements, in order — stop at the step that convicts the cause, and don't buy RAM until step 5.5 entries
Symptom · 01
Executors dying mid-shuffle with OutOfMemoryError or exit 143 and you don't know if it's skew or size
→
Fix
Open the stage's task table, sort by shuffle-read bytes descending, and divide the top task by the median. Over 10x means skew — skip to salting and don't touch memory settings. Near 1x with all tasks fat means genuine capacity pressure — proceed to partition math. This single ratio decides the whole response, and it takes 60 seconds in the Spark UI.
Symptom · 02
Confirmed skew: a handful of tasks hold gigabytes while the rest hold megabytes
→
Fix
Find the hot key with a GROUP BY key ORDER BY COUNT(*) DESC LIMIT 20 on the shuffle input — if the top key owns over ~10% of rows, salt it: add a random 0-31 suffix on the large side, replicate the small side across all 32 suffixes, join, then strip the suffix and re-aggregate. For Spark 3 joins, try the SKEW hint or AQE skew splitting first (spark.sql.adaptive.skewJoin.enabled, on by default) — they automate this for supported join types.
Symptom · 03
Uniform fat tasks: every partition is 1-2GB and executors die evenly across the stage
→
Fix
Do partition math: total shuffle bytes divided by 200MB target size equals the partition count you need. Set spark.sql.shuffle.partitions (default 200 — almost always too low past a few hundred GB) or call repartition(N) explicitly, and verify with the bucket histogram query that p99 partition size lands 128-256MB. Recheck the task table after one test run before changing any memory setting.
Symptom · 04
OOM persists at healthy partition sizes with high GC time (over ~10% of task time)
→
Fix
Attack object overhead, not bytes: switch spark.serializer to Kryo (spark.serializer=org.apache.spark.serializer.KryoSerializer), replace Python UDFs with SQL built-ins to cut deserialization churn, and unpersist cached DataFrames the moment their last action completes — retained cache silently steals unified memory from execution. Enable spark.memory.offHeap.enabled with spark.memory.offHeap.size (e.g. 4-8G) so Tungsten sorts spill off-heap instead of into the GC's path.
Symptom · 05
Healthy partitions, lean code, freed caches — and tasks still die near the heap ceiling
→
Fix
Now, and only now, size the box: raise spark.executor.memory one step (e.g. 8G→12G) and add matching memoryOverhead (max(384M, 10%) is thin for shuffle-heavy stages — set spark.executor.memoryOverhead explicitly to 2-4G). Watch YARN's kill message: if it still names physical-memory over-limit, overhead was the gap; if heap fills flat to the ceiling, add cores' worth of parallelism instead — more, smaller tasks beat fewer, bigger boxes.
Shuffle OOM causes compared
Root CauseHow to ConfirmFixPrevention
Skewed hot key in shuffleTop task over 10x median bytes; key census shows one key above ~10% of rowsSalt N = hot bytes / 200MB; SKEW hint or AQE split for moderate skewKey-share gates in CI; AQE skew handling left enabled
Too few shuffle partitionsTasks uniformly fat (p99 near p50, both over ~1GB); default 200 partitionsSet partitions to shuffle bytes / 200MB; verify p99 lands 128-256MBPartition math in design review; per-stage counts, not global hope
Serialization and GC overheadGC time above ~10%; Python RSS leads JVM heap; Java serializer in confKryo serializer; SQL built-ins over UDFs; off-heap sort workspaceSerializer lint in submit templates; UDF budget per code review
Retained cache stealing executionStorage tab shows stale full caches during shuffle; heap flat at ceilingUnpersist after last dependent action; MEMORY_AND_DISK over MEMORY_ONLYCache lifecycle audit per job; alert on stale cached GB during shuffles
⚙ Quick Reference
5 commands from this guide
FileCommand / CodePurpose
unified_memory_math.pyheap_mb = 8 * 1024Unified Memory
skew_vs_size_census.sqlSELECT device_id, COUNT(*) AS n,Skew vs Genuine Size
salted_skew_join.pyfrom pyspark.sql import functions as FSalting Hot Keys
PartitionMath.scalaval shuffleBytes = 9.4 * 1024L * 1024L * 1024L * 1024L // 9.4 TBPartition Counts
lean_shuffle_conf.pyspark.conf.set("spark.serializer",Off-Heap, Kryo, and Cache Hygiene

Key takeaways

1
Shuffle OOMs are demand problems first
measure the unified ceiling (~60% of heap minus 300MB) and fit partitions to it.
2
Top-to-median task bytes above ~10x convicts skew
salt with N = hot bytes / 200MB, don't buy RAM.
3
Uniform fat tasks convict the 200-partition default
divide shuffle bytes by 200MB and verify p99 at 128-256MB.
4
Kryo, off-heap sorts, and SQL built-ins shrink per-byte pressure that partition math alone can't reach.
5
Every persist() is a loan
unpersist after the last dependent action or stale cache starves later shuffles.
6
Size the box last and step-wise
overhead for physical kills, heap for Java fills, fewer cores for cheap headroom.

Common mistakes to avoid

5 patterns
×

Doubling executor memory fleet-wide for a skewed stage

Symptom
Bill doubles on every run while one immortal task dies later on a bigger box — 32G executors dying at hour 7 instead of 8G at hour 4
Fix
Read the top-to-median task ratio first. Above ~10x means salt the hot key; memory only helps uniform-fat stages.
×

Tuning spark.memory.fraction and storageFraction blindly

Symptom
Fractions shuffled for days with identical OOMs — the unified ceiling moved 10% while the hot partition overshoots it 10x
Fix
Compute the unified ceiling once, then fix demand (partition size, skew) to fit it. Touch fractions only for borderline cases within ~20% of the ceiling.
×

Leaving spark.sql.shuffle.partitions at the 200 default on terabyte shuffles

Symptom
Uniform 40GB+ tasks dying evenly across the stage — every executor genuinely asked to swallow a door-sized slab
Fix
Divide shuffle bytes by 200MB and set the count explicitly per big stage; verify p99 task size lands 128-256MB.
×

Caching inputs and forgetting them across a multi-stage job

Symptom
Storage tab shows 100%-cached frames no action needs while shuffle tasks starve — a slow leak that kills stage 9 for stage 2's convenience
Fix
Unpersist immediately after the last dependent action; prefer MEMORY_AND_DISK so spill absorbs spikes instead of eviction storms.
×

Serializing shuffles with default Java serialization plus Python UDFs

Symptom
GC time above 30%, Python RSS dwarfing JVM heap, shuffle bytes 2-10x larger than the same data in Parquet
Fix
Switch to Kryo, replace UDFs with null-safe SQL built-ins, and move sort workspace off-heap before concluding the box is small.
INTERVIEW PREP · PRACTICE MODE

Interview Questions on This Topic

Q01JUNIOR
An executor dies with OutOfMemoryError mid-shuffle. What's your first me...
Q02SENIOR
Explain Spark's unified memory region and the eviction asymmetry.
Q03SENIOR
How do you salt a skewed join, and how do you pick N?
Q04SENIOR
When do you choose memoryOverhead over heap, and when fewer cores over b...
Q05SENIOR
AQE skew splitting is enabled but a 1TB hot key still kills the task. Wh...
Q01 of 05JUNIOR

An executor dies with OutOfMemoryError mid-shuffle. What's your first measurement?

ANSWER
Top-to-median shuffle-read bytes in the stage's task table. Over ~10x means skew (salt); near 1x and fat means sizing (partitions). That ratio decides the whole response before any setting changes.
FAQ · 6 QUESTIONS

Frequently Asked Questions

01
What are spark.memory.fraction and storageFraction defaults?
02
How big should shuffle partitions be?
03
What is salting and when do I need it?
04
Does AQE handle skew automatically?
05
When does off-heap memory help shuffle OOMs?
06
Why did doubling memory make my job slower instead of fixing it?
N
Naren Founder & Principal Engineer

20+ years shipping production backend systems. Drawn from code that ran under real load.

Follow
✓ Verified
production tested
September 27, 2026
last updated
2,085
articles · all by Naren
🔥

That's Spark. Mark it forged?

6 min read · try the examples if you haven't

←
Previous
Spark Job Aborted: Task Failed 4 Times — Read the Real Cause
2 / 5 · Spark
Next
Spark Data Skew: One Task Runs for Hours
→