spark shuffle out of memory partition sizing
Fixes Spark shuffle out-of-memory errors through partition sizing. Use when stages fail with executor OOM during shuffles, when you see huge spill or single giant tasks, or when the default 200 shuffle partitions are wrong for your data size. Not for AnalysisException column resolution, for to_timestamp parse failures, or for broadcast join threshold tuning.
TL;DR
Set spark.sql.shuffle.partitions to roughly 2-3x your core count (or size partitions to 100-200MB each), and repartition skewed keys before the shuffle. Most shuffle OOMs are a handful of giant partitions, not a cluster-wide shortage.
spark shuffle out of memory partition sizingUse this when
- Stages fail with executor lost or OOM during a shuffle (join, groupBy, repartition)
- The Spark UI shows a few tasks taking far longer than the rest
- Spill to disk is huge on shuffle stages
Not for this skill when
- Analysis fails to resolve a column name
to_timestampproduces nulls on malformed rows- You are deciding whether a join should broadcast
Steps
- Find the skew in the Spark UI or from the data. The symptom is task duration imbalance:
df.groupBy("key").count().orderBy("count", ascending=False).show(10)Expected output: the key distribution. If the top key holds a large share of rows, no partition count fixes it, you need skew handling (step 4).
- Size shuffle partitions so each holds 100-200MB:
spark.conf.set("spark.sql.shuffle.partitions", "400")Expected output: set before the query runs. Rule of thumb: total shuffle bytes divided by 128MB, rounded up, and at least 2-3x the number of cores so tasks stay parallel. The default 200 is wrong for both tiny and huge shuffles.
- For a single heavy DataFrame, repartition explicitly on the join/group key:
df = df.repartition(400, "key")Expected output: the data spread evenly across 400 partitions before the expensive operation. This beats the global setting when one stage dominates the job.
- Handle skewed keys with salting when one key dominates:
from pyspark.sql import functions as F
salted = df.withColumn("salt", (F.rand() * 10).cast("int"))
agg = salted.groupBy("key", "salt").agg(F.sum("v").alias("v")) \
.groupBy("key").agg(F.sum("v").alias("v"))Expected output: the hot key spreads across 10 salt buckets for the partial agg, then the small partial results combine. Two-stage aggregation keeps any single task small.
- Verify in the UI: after the fix, task durations should be roughly even and spill near zero:
# in the Spark UI Stages tab: check max vs median task duration and shuffle spillExpected output: max task duration close to the median. If one task still dominates, the skew is on a different key than you salted.
Variant phrasings
spark executor lost shuffle fetch failed
Often the downstream symptom of OOM during shuffle. Fix the partition sizing first; if fetches still fail, look at network and external shuffle service health.
spark java heap space groupBy
The groupBy shuffle produced partitions bigger than executor memory. Reduce per-partition bytes via steps 2-4.
spark.sql.shuffle.partitions best value
There is no universal value: aim for 100-200MB per partition. For a 100GB shuffle on 50 cores, 800-1000 partitions is a sane start.
Why it happens
A shuffle redistributes rows by key across partitions, and each partition must fit in executor memory during the reduce side. Too few partitions makes each one huge; skewed keys make one partition huge regardless of count. The defaults were chosen for modest data and rarely survive contact with production volumes.
Edge cases
- Adaptive Query Execution (AQE) can coalesce small partitions automatically at runtime, keep
spark.sql.adaptive.enabled=true. - Over-partitioning has its own cost: thousands of tiny tasks add scheduling overhead and tiny files on write.
- Sort-merge joins shuffle both sides, size partitions for the larger side.
repartitiontriggers a full shuffle itself, usecoalescewhen only reducing partition count without a shuffle.
Provenance
Resolved from the public thread: https://vectle.com/posts/pst_unDQSMirFGmAc-YrTLjUOw
Maintainer review
No maintainer verification is recorded for this version.
This records the version a maintainer checked. It does not assert that the version is the latest upstream release.