Spark Data Skew: One Task Runs While 199 Idle
One task runs for hours while 199 sit idle: that's data skew.
20+ years shipping production backend systems. Everything here is grounded in real deployments.
- ✓Spark stages, tasks, and hash partitioning basics
- ✓GROUP BY aggregations and join types in Spark SQL
- ✓Reading task-duration histograms in the Spark UI
- Data skew means a few keys own most rows, so hash partitioning piles gigabytes onto one task while hundreds finish in minutes — read the task-duration histogram first
- Confirm with a GROUP BY key census: any key above ~10% of rows is skewed, and top-to-median task bytes above ~10x convicts it beyond doubt
- Salt the hot key (random 0-31 suffix on the big side, replicate the small side), or use the SKEW join hint and AQE skew splitting on Spark 3
- Prevent repeats with two-stage aggregation (partial plus final), hashed-key repartitions, and key-share gates in CI that fail builds before skewed joins ship
Imagine a post office where 200 clerks sort by zip code — except one zip holds a warehouse sending 11% of all parcels. One clerk drowns while 199 finish in minutes. That's skew: Spark deals work by key, so one popular key means one giant task. Fixes all split the pile: deal the warehouse mail across clerks (salting), teach the sorter to split overloaded bins (AQE), or pre-sort so no bin overflows (hashed repartitioning).
Every Spark engineer meets skew the same way: a job that's 99% done for three hours. The UI shows 199 tasks green in minutes and one task still running — then dying, then retrying, then dying again. The progress bar is technically honest (99% of tasks finished) and completely misleading (the remaining 1% holds 11% of your data). Newcomers blame the cluster. Veterans check the keys.
Skew is a data-shape problem wearing a performance costume. Hash partitioning assumes keys spread evenly; real data has celebrities — a null device_id from a bot farm, a default 'UNKNOWN' customer, a viral product ID. Every row with the same key lands on the same partition, and one partition grows until no executor can swallow it. Memory tuning can't fix it (the task is atomic), retries can't fix it (every attempt chokes identically), and bigger boxes barely dent it.
This article covers the full skew playbook in the order professionals apply it. Diagnose with histograms and key censuses, fix hot joins with salting and SKEW hints, let Adaptive Query Execution split what it can, restructure aggregations into two stages, and build CI gates so the next celebrity key fails a test instead of a 2 AM job.
Reading the Histogram: Skew's Mugshot
Skew announces itself in the task table before anywhere else: durations like 4 min, 5 min, 4 min ... 7 hours. That shape — a flat plain of fast tasks with one Everest — is pathognomonic. No hardware fault produces it (sick disks slow tasks 2-5x, not 100x), no sizing fault produces it (uniform fat tasks make a flat plain of slow tasks), and no code bug produces it (poison data kills tasks fast, it doesn't stretch them). When you see Everest, you have skew, and every minute spent on memory settings is wasted.
Two ratios sharpen the mugshot. Duration ratio (longest task / median task) above ~10x convicts skew on time; shuffle-read ratio (fattest task / median task) above ~10x convicts it on bytes. Both high means a giant partition doing giant work; duration high with bytes even means one task is stuck (GC death spiral or spilled-sort thrash on merely-large data) — still skew-adjacent, still fixed by splitting, but the census matters more. Screenshot both ratios into the incident channel: they end the 'add memory' debate instantly because no box size bridges a 30x gap.
The timeline view adds the alibi check. Skew's Everest starts slow from the beginning (it reads its mountain from second one), while a straggler-disk task runs normally then flatlines. If the longest task's read curve climbs steadily for hours, that's a giant partition being consumed — salt it. If it flatlines mid-task, that's thrash — still split it, and check spill metrics to confirm.
The Key Census: Naming the Celebrity
The histogram proves skew; the census names the key. GROUP BY the join or aggregate key over the shuffle input, order by count descending, and read the top rows like a most-wanted list. In production the winners are depressingly familiar: NULL (bot traffic, missing instrumentation), empty string, 0, 'UNKNOWN', a test account, a viral product, a NAT-gateway device_id. Real identifiers almost never skew — placeholders and shared infrastructure do.
Read shares, not just counts: a top key at 13% of rows (our NULL) needs structural treatment (separate lane or heavy salting), while a top key at 2% might yield to AQE splitting alone. Also check the tail shape — one celebrity versus a heavy-hitter band of thousands of warm keys changes the weapon: single celebrities get salted or laned, bands get broadcast hints (if the small side qualifies under the 10MB auto-broadcast threshold) or two-stage aggregation.
Run the census on a sample when inputs are huge (TABLESAMPLE or a single recent partition usually reveals celebrities — skew is fractal, the hot key is hot in every slice), but always census the exact input of the dying stage, not an upstream approximation. Filters between your census and the stage can create or destroy skew: a filter that removes 90% of ordinary rows while keeping all NULLs concentrates the celebrity from 1% to 10% silently.
Salting: Splitting One Key Across N Tasks
Salting is controlled demolition for hot keys. You transform key K into N sub-keys K_0..K_{N-1} by appending a random suffix on the large side, replicate each small-side row N times (once per suffix), join on the salted key, then strip the suffix and re-aggregate to restore the original grain. The hot partition's bytes divide by N (modulo randomness), and N ordinary tasks replace one immortal task. With N=16, our 830GB NULL partition would have become 16 tasks near 52GB — still large, which is why NULLs deserved a separate lane instead; salting shines on real keys (viral products, gateway IDs) that must participate in the join.
Size N from the hot key's bytes: N = ceil(hot_key_bytes / 200MB). Undershoot and sub-tasks still OOM; overshoot and small-side replication (N copies of every row) bloats the shuffle. Powers of two (16, 32, 64) keep the mental math clean. Salt only the skewed side when the other side is small — and if the small side fits under the broadcast threshold (10MB default), skip salting entirely and broadcast it: a broadcast join has no shuffle, hence no skew, full stop.
The re-aggregation step is where salting silently corrupts results if rushed. Counts and sums re-aggregate safely (sum the partial sums); averages don't (average the partial sums and you're wrong — carry sum and count, divide at the end); distinct counts need partial HLL sketches or a two-pass exact approach. Write the grain-restoration as its own reviewed query, and assert row counts before and after: salted output grouped back to the original key must match the unsalted logical result on a skew-free sample.
SKEW Hints and AQE Splitting: The Automated Net
Spark 3 gives you two automated skew fighters before hand-salting. The SKEW hint (SELECT /+ SKEW('table_name', 'column') / ...) declares which relation and column skew, letting the optimizer split those shuffle partitions and replicate the other side's matching rows — essentially compiler-assisted salting without rewriting your query. It works on joins and some aggregations, and its virtue is locality: the skew knowledge lives in the query text where reviewers see it, not in a 40-line salting preamble.
Adaptive Query Execution's skew handling (spark.sql.adaptive.enabled plus spark.sql.adaptive.skewJoin.enabled, both on by default since Spark 3.2) goes further by detecting skew at runtime: it measures actual shuffle partition sizes mid-query and splits oversized ones (defaults: skewedPartitionThresholdInBytes 256MB, skewedPartitionFactor 5), replicating the corresponding small-side partitions. No hints, no rewrite — the plan adapts to the data it actually sees, which also covers skew you never predicted.
Know their ceilings. Hints only help supported operators (skewed aggregations without joins may ignore them), and AQE's defaults split modestly — a single key worth 11% of a 9TB shuffle defeats factor-5 splitting the way a river defeats a sandbag. Treat automation as the first attempt for moderate skew (top key under ~5%) and as a permanent second net under hand-salting for extreme skew: even perfect salt benefits from AQE catching the next key nobody censused.
Two-Stage Aggregation and Hashed Repartitioning
Skewed aggregations (GROUP BY on a hot key, no join involved) can't borrow join tricks — there's no small side to replicate. The answer is two-stage aggregation: first group by (key, salt) with a partial aggregate, then group the partials by key into the final result. Stage one spreads the hot key across N tasks; stage two combines N partial rows per key — trivial work, since N rows fit anywhere. This is map-side-combine generalized: you're paying one extra shuffle to convert an atomic mountain into N pebbles plus one handful.
Hashed-key repartitioning prevents skew rather than curing it. Writing upstream outputs with repartition(N, hash(key)) (or bucketed tables keyed on the join column) guarantees downstream stages inherit balanced partitions — the skew never forms because no single partition accumulates a celebrity. Size N from output bytes / 200MB as usual, and prefer it at pipeline boundaries: one balanced write protects every downstream reader, while per-query salting must be re-applied at each use.
Broadcast discipline completes the set. Any dimension under the auto-broadcast threshold (spark.sql.autoBroadcastJoinThreshold, 10MB default — raise deliberately to 50-100MB when driver memory allows) should broadcast, deleting the shuffle and the skew surface together. Check the SQL tab's plan after every supposedly-broadcast join: type mismatches between join keys silently disable broadcasting (int vs string never broadcasts), and that silent fallback is itself a classic skew factory.
Making Skew Unshippable: Gates, Alerts, Habits
Skew is a data-distribution property, so it returns whenever data drifts — today's balanced key is next quarter's celebrity after a product launch, a bot surge, or an upstream default change. The durable fix is cultural: make skewed shapes fail in CI instead of in production. A census job that samples each new join's inputs and fails the merge when any key exceeds ~5% of sampled rows costs minutes per PR and saves 2 AM pages per quarter.
Production alerts should watch ratios, not just durations. Max-to-median task time per stage above ~8x, or any task exceeding 5x its stage's p50 two runs consecutively, pages the owning team with the stage ID and the histogram attached. Duration alerts alone fire too late (the job already burned hours); ratio alerts fire while the Everest is still climbing, when killing and salting saves the run instead of the postmortem.
Habits close the loop. Every new join gets its top-20 key census pasted into the PR description. Every pipeline boundary write uses hashed-key repartitioning sized from bytes. AQE skew handling stays enabled cluster-wide (disabling it 'for determinism' deletes your free second net). And the incident ritual stays fixed: histogram first, census second, memory never — until the ratios say sizing, which, once these gates hold, they rarely do.
The NULL Key That Ate Black Friday: 7 Hours for One Task
- NULL is the most common celebrity key in production and the easiest to miss — hash partitioners route every NULL identically, so a bot surge becomes a single-partition bomb. Census NULL share explicitly on every join key, not just distinct values.
- Speculation on a skewed stage doubles the damage: the duplicate reads the same giant partition and dies the same death. Check the top-to-median ratio before enabling speculation — above ~10x, speculate never.
- Separate lanes beat salting for semantically separable hot keys. Bots don't purchase, so their rows never needed the join at all — one filter outperformed any N-way replication on cost, clarity, and runtime.
| File | Command / Code | Purpose |
|---|---|---|
| key_census.sql | SELECT device_id, COUNT(*) AS n, | The Key Census |
| salted_join.py | from pyspark.sql import functions as F | Salting |
| skew_hint_aqe.sql | SELECT /*+ SKEW('events', 'device_id') */ | SKEW Hints and AQE Splitting |
| two_stage_agg.py | from pyspark.sql import functions as F | Two-Stage Aggregation and Hashed Repartitioning |
Key takeaways
Common mistakes to avoid
5 patternsEnabling speculation on a skewed stage
Raising shuffle partitions to 'fix' extreme single-key skew
Averaging partial averages in two-stage aggregation
Salting both sides symmetrically when only one side skews
Forgetting NULLs in the key census
Interview Questions on This Topic
A job sits at 199/200 tasks for 3 hours. What's your first hypothesis and first query?
Frequently Asked Questions
20+ years shipping production backend systems. Everything here is grounded in real deployments.
That's Spark. Mark it forged?
6 min read · try the examples if you haven't