All concepts

Data Skew

199 tasks finish in a minute and one runs for two hours — the cluster isn't slow, one key owns half your data.

Distributed Processing · Advanced · ~5 min

In plain English

Ten checkouts open and one queue has half the shop in it. Opening an eleventh checkout doesn't help the people already in that queue.

Why it's worth your time

It's the reason jobs get slower as data grows even though the cluster got bigger.

If you remember three things

  • One task can't be split, so more executors change nothing
  • Diagnose by max ÷ median task duration, not by absolute time
  • Null and sentinel keys are the cheapest skew to remove

Overview

Distributed work is only as fast as its slowest task, and skew is what makes one task enormously slower than the rest. It happens whenever the partitioning key is unevenly distributed: a guest_user id used by every anonymous session, a null that stands in for missing, a single tenant a hundred times bigger than the rest, a default timestamp. Adding executors does nothing, because the problem isn't capacity — one task cannot be split across machines. The diagnosis is always the same: compare the maximum task duration and shuffle-read size against the median.

In an interview

Skew is uneven key distribution, so one partition holds far more rows than the others and its task dominates the runtime. Adding resources doesn't help — a single task can't be parallelised. Diagnose by comparing max task duration and shuffle-read bytes against the median; fix by salting the hot key, isolating and handling it separately, or letting adaptive execution split the partition.

Production defaults

First check
group by the key and look at the top 5 counts
Free fix
filter nulls and sentinels before any wide operation
Salting
size N from the actual imbalance; large N replicates the small side pointlessly

What breaks

  • 199 tasks done, 1 running for hours — Textbook skew. Find the dominant key; filter it if spurious, salt or isolate it if real.
  • Salted groupBy gives wrong totals — Missing the second aggregation pass that combines across salts.

Watch it explained

Fix Spark Joins Getting Stuck at 99%! | Handle Data Skew in PySpark with Salting — Sriw World of Coding, 6:13

Related