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

```text
spark shuffle out of memory partition sizing
```

## Use 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_timestamp` produces nulls on malformed rows
- You are deciding whether a join should broadcast

## Steps

1. Find the skew in the Spark UI or from the data. The symptom is task duration imbalance:

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

2. Size shuffle partitions so each holds 100-200MB:

```python
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.

3. For a single heavy DataFrame, repartition explicitly on the join/group key:

```python
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.

4. Handle skewed keys with salting when one key dominates:

```python
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.

5. Verify in the UI: after the fix, task durations should be roughly even and spill near zero:

```python
# in the Spark UI Stages tab: check max vs median task duration and shuffle spill
```
Expected 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.
- `repartition` triggers a full shuffle itself, use `coalesce` when only reducing partition count without a shuffle.

## Provenance

Resolved from the public thread: https://vectle.com/posts/pst_unDQSMirFGmAc-YrTLjUOw
