Home › Data Engineering › Spark Data Skew: One Task Runs While 199 Idle
Advanced 6 min · September 23, 2026

Spark Data Skew: One Task Runs While 199 Idle

One task runs for hours while 199 sit idle: that's data skew.

N
Naren Founder & Principal Engineer

20+ years shipping production backend systems. Everything here is grounded in real deployments.

Follow
✓ Production
production tested
September 27, 2026
last updated
2,085
articles · all by Naren
Before you start⏱ 15 min
  • ✓Spark stages, tasks, and hash partitioning basics
  • ✓GROUP BY aggregations and join types in Spark SQL
  • ✓Reading task-duration histograms in the Spark UI
 ● Production Incident 🔎 Debug Guide
⚡Quick Answer
  • 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
✦ Definition~90s read
What is Spark Data Skew?

Data skew is uneven key distribution meeting hash partitioning's uniformity assumption. Spark assigns rows to partitions with hash(key) mod N — perfectly even when keys spread evenly, catastrophic when one key owns 13% of rows. Every row sharing that key lands on the same partition, processed by one task, on one executor, with one heap.

★
Imagine a post office where 200 clerks sort by zip code — except one zip holds a warehouse sending 11% of all parcels.

The other N-1 tasks finish in minutes; the celebrity task runs for hours, spills constantly, GC-thrashes, and usually dies. Retries re-execute the identical mountain; speculation duplicates it; memory upgrades grow the box around a task that's atomic by construction.

Skew hides in ordinary-looking data because the skewing values are usually semantically empty: NULLs from missing instrumentation, 'UNKNOWN' defaults, zero IDs, test accounts, shared gateway identifiers. Distinct counts stay healthy (millions of distinct keys!) while one value quietly owns a tenth of all rows.

Filters can manufacture it mid-pipeline — removing 90% of normal rows while keeping every NULL concentrates a 1% curiosity into a 10% killer without any upstream change.

The fix family shares one idea: break the atomicity. Salting divides one key into N sub-keys, two-stage aggregation divides one combine into N partials plus a trivial final, separate lanes remove non-participating rows from the join entirely, and AQE splitting divides oversized shuffle partitions at runtime.

All four convert one impossible task into many ordinary ones — the only transformation skew respects.

Plain-English First

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.

📊 Production Insight
Teams that screenshot the two ratios into the incident channel within 10 minutes resolve skew pages in under an hour; teams that start with memory bumps average 3+ days because each test run costs a full job duration to learn nothing.
🎯 Key Takeaway
A flat plain of fast tasks plus one Everest is skew's unmistakable mugshot — confirm with longest-to-median ratios on duration and bytes before touching anything else.

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.

key_census.sqlSQL
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
-- Name the celebrity key: run against the dying stage's exact input.
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;
-- Verdict guide: top key > 10% = extreme (salt or separate lane),
-- 2-10% = moderate (SKEW hint or AQE split), < 2% = not skew.

-- NULL share deserves its own explicit line: hash partitioning
-- routes every NULL to one partition deterministically.
SELECT COUNT_IF(device_id IS NULL) AS nulls,
       COUNT(*) AS total,
       ROUND(100.0 * COUNT_IF(device_id IS NULL) / COUNT(*), 2) AS null_pct
FROM lake.events
WHERE dt = '2026-09-22';
📊 Production Insight
A census habit (top-20 keys on every new join, pasted into the PR) catches NULL and 'UNKNOWN' celebrities before merge — one team blocked 14 skewed joins in code review the quarter after adopting it.
🎯 Key Takeaway
Census the dying stage's exact input and read shares: NULLs and placeholders are the usual celebrities, and the top share (above or below ~10%) picks the weapon.

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.

salted_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
29
30
from pyspark.sql import functions as F

N = 16  # hot key ~3.2GB -> 16 x ~200MB sub-tasks

big = spark.read.parquet("s3://lake/events/dt=2026-09-22/")
# Filter semantic non-participants to their own lane BEFORE salting.
null_lane = big.filter(F.col("device_id").isNull())
big_salted = (
    big.filter(F.col("device_id").isNotNull())
    .withColumn("salted_key",
                F.concat(F.col("device_id"), F.lit("_"),
                         (F.rand() * N).cast("int")))
)

small = spark.read.parquet("s3://lake/devices/")
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")
# Restore grain: strip suffix, carry sum AND count (never avg partials).
restored = (
    joined.withColumn("device_id", F.split("salted_key", "_")[0])
    .groupBy("device_id")
    .agg(F.sum("amount").alias("total"), F.count("*").alias("n"))
    .withColumn("avg", F.col("total") / F.col("n"))
)
restored.write.mode("overwrite").parquet("s3://lake/joined/")
📊 Production Insight
The re-aggregation step is where a rushed salt corrupts revenue numbers: one team averaged partial averages for a quarter before an audit caught it — carry sum and count, divide once, and assert grains match on a clean sample.
🎯 Key Takeaway
Salt as suffix-plus-replicate with N from hot-bytes/200MB, lane out non-participants first, and restore grain with sum-and-count — never average partial averages.

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.

skew_hint_aqe.sqlSQL
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
-- Declare the skew: optimizer splits these partitions for you.
SELECT /*+ SKEW('events', 'device_id') */
  e.device_id, d.model, SUM(e.amount) AS total
FROM lake.events e
JOIN lake.devices d ON e.device_id = d.device_id
WHERE e.dt = '2026-09-22'
GROUP BY e.device_id, d.model;

-- Runtime second net (defaults since Spark 3.2; verify, don't assume):
-- spark.sql.adaptive.enabled = true
-- spark.sql.adaptive.skewJoin.enabled = true
-- spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes = 256MB
-- spark.sql.adaptive.skewJoin.skewedPartitionFactor = 5
-- For extreme single-key skew, raise the factor before hand-salting:
-- SET spark.sql.adaptive.skewJoin.skewedPartitionFactor = 10;
📊 Production Insight
AQE splitting silently fixed 70% of one fleet's moderate skews after the 3.2 upgrade — the remaining 30% (single keys above ~8% of shuffle) all needed hand-salting, and the hint-plus-AQE combo meant salt code only shipped where automation provably failed.
🎯 Key Takeaway
Declare skew with SKEW hints and keep AQE splitting enabled as the runtime net — then hand-salt only the extreme cases automation provably can't split.

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.

two_stage_agg.pyPYTHON
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
from pyspark.sql import functions as F

N = 32
raw = spark.read.parquet("s3://lake/events/dt=2026-09-22/")

# Stage 1: partial aggregate spread across (key, salt) — N tasks share the load.
partial = (
    raw.withColumn("salt", (F.rand() * N).cast("int"))
    .groupBy("device_id", "salt")
    .agg(F.sum("amount").alias("part_sum"), F.count("*").alias("part_n"))
)

# Stage 2: N partial rows per key combine anywhere — trivially small.
final = (
    partial.groupBy("device_id")
    .agg(F.sum("part_sum").alias("total"), F.sum("part_n").alias("n"))
    .withColumn("avg", F.col("total") / F.col("n"))
)

# Balanced writes protect downstream readers: hash-partition the output.
final.repartition(512, "device_id").write.mode("overwrite").parquet(
    "s3://lake/device_totals/")
💡Carry sum and count through every partial stage
Partial averages, partial distinct-counts, and partial medians don't compose — averaging averages is wrong by construction. Carry additive quantities (sum, count, HLL sketches) through stage one and compute the delicate statistic once in stage two. Assert final grain matches a clean-sample baseline before shipping.
📊 Production Insight
Two-stage aggregation fixed a skewed COUNT DISTINCT in one pass where salting couldn't apply (no join existed): partial HLL sketches per (key, salt) merged exactly in stage two, cutting a 5-hour Everest to 22 minutes.
🎯 Key Takeaway
No join means no salt target — use two-stage (key, salt) partials plus hashed repartitioned writes, and broadcast every dimension that fits.

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.

📊 Production Insight
One quarter after adopting census-in-PR plus ratio alerting, a 60-pipeline fleet went from 11 skew pages to zero — the gates caught 14 celebrity keys in review, and the two that escaped were auto-split by AQE before anyone noticed.
🎯 Key Takeaway
Skew recurs with data drift, so gate it in CI (key share above ~5% fails the merge) and alert on task ratios in prod — culture, not config, is the permanent fix.
● Production incidentPOST-MORTEMseverity: high

The NULL Key That Ate Black Friday: 7 Hours for One Task

Symptom
The click-to-purchase join (6.4 TB shuffle, 200 partitions) finished 199 tasks in under 5 minutes each night of Thanksgiving week, then sat at 199/200 for hours. The last task's shuffle-read climbed past 800GB, GC time passed 60%, and it died with heap exhaustion at hour 7 — on both the scheduled run and the manual re-run. Executives watched a stuck-at-99% dashboard through all of Black Friday morning while revenue numbers stayed frozen.
Assumption
The team assumed a straggler disk: one slow node holding up an otherwise healthy stage. They enabled speculation, then blacklisted the host the long task ran on, then raised shuffle partitions 200→2,000. Each change moved the long task to a new host where it died identically. Speculation actually doubled the damage — the duplicate copy read the same 800GB and died too, burning twice the I/O for zero information.
Root cause
A key census (run in minute 10 of the real triage, on day 3) showed NULL device_id on 840M of 6.4B rows — 13% of the shuffle. Bot traffic without device identifiers had tripled that week, and Spark's hash partitioner sends every NULL to the same partition deterministically. One task owned 830GB while the median owned 28GB (a 30x ratio). No disk, host, or memory setting was ever involved: any task handed 830GB on an 8G executor dies, on any host, at any partition count below ~5,000 — and even 5,000 leaves a 166GB giant.
Fix
NULLs got their own lane: bot traffic (device_id IS NULL) was filtered into a separate aggregate that doesn't join at all (bots don't purchase), cutting the hot key to zero rows. The remaining join got 16-way salting on the top 50 keys plus AQE skew handling left enabled as a second net, and shuffle partitions went to 8,192 for ~800MB... no — verified at ~200MB p99 after the NULL lane removal shrank the shuffle to 5.5 TB. Runtime went from 7 dying hours to 61 green minutes. A CI gate now fails any merge whose key census shows a single key above 5% of sampled rows.
Key lesson
  • 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.
Production debug guideEach step either convicts skew or eliminates it — stop tuning memory the moment step 1 says skew.5 entries
Symptom · 01
Job stuck at 99% with one or two tasks running hours longer than the rest
→
Fix
Open the stage's task table, sort by duration descending, and compute longest-task versus median-task time and shuffle bytes. Ratio above ~10x on both axes is skew — screenshot it for the incident channel and skip all memory tuning. Below ~3x with everything slow is a sizing or code problem instead; follow the executor-OOM playbook for that family.
Symptom · 02
Histogram convicts skew but you don't know which key is hot
→
Fix
Run SELECT key, COUNT() AS n, 100.0COUNT()/SUM(COUNT()) OVER () AS pct FROM input GROUP BY key ORDER BY n DESC LIMIT 20 on the shuffle input. Any key above ~10% is your celebrity; NULL or defaults ('UNKNOWN', '', 0) in the top rows are the classic culprits. Note the exact share — it sets the salt factor N = hot bytes / 200MB.
Symptom · 03
Hot key found in a join, Spark 3+, moderate skew (top key under ~5% of rows)
→
Fix
Try the cheap automation first: SELECT /+ SKEW('events', 'device_id') / ... to declare the skew, and verify AQE skew handling is on (spark.sql.adaptive.enabled and skewJoin.enabled both default true since 3.2). Re-run and check whether the longest task dropped under 2x median. If it did, ship it with a comment naming the key — if not, the skew is extreme and needs hand-salting.
Symptom · 04
Extreme skew (single key above ~10%) or AQE hints didn't move the longest task
→
Fix
Hand-salt: suffix the big side's key with a random 0..N-1, cross-join the small side against N suffix rows and concatenate the same suffix, join on the salted key, then strip the suffix and re-aggregate to restore grain. Start N at hot-key-bytes/200MB (16-64 typical). For aggregations (no join), use two-stage group-by: partial aggregate by (key, salt), then final aggregate by key.
Symptom · 05
Skew fixed today but you fear next quarter's celebrity key
→
Fix
Make skew unshippable: add a CI census job that samples join inputs and fails the merge when any key exceeds 5% of rows, repartition high-cardinality writes by hashed key (repartition(N, col) with N from bytes/200MB) so downstream stages inherit balance, and keep AQE skew splitting enabled as the runtime second net. Alert on max-to-median task ratio per stage in production.
Skew responses compared
Root CauseHow to ConfirmFixPrevention
Single celebrity key in a join (NULL, default, viral ID)Census shows one key above ~10%; task ratio above ~10xSalt N = hot bytes / 200MB, or separate lane if rows needn't joinCI key-share gate at 5%; explicit NULL lanes in contracts
Moderate skew (top key 2-10%)Census top key 2-10%; Everest 3-10x medianSKEW hint plus AQE skew splitting; tune factor for stubborn casesLeave AQE skew handling enabled cluster-wide as second net
Skewed aggregation with no joinGROUP BY stage shows Everest; no small side exists to replicateTwo-stage (key, salt) partials; carry sum/count, compute delicates oncePre-aggregate upstream; hashed repartition on hot group keys
Broadcast fallback silently disabledSQL plan shows SortMergeJoin where broadcast expected; key types differFix key types (int vs string); raise broadcast threshold deliberatelyPlan assertion in tests: explain must show BroadcastHashJoin
⚙ Quick Reference
4 commands from this guide
FileCommand / CodePurpose
key_census.sqlSELECT device_id, COUNT(*) AS n,The Key Census
salted_join.pyfrom pyspark.sql import functions as FSalting
skew_hint_aqe.sqlSELECT /*+ SKEW('events', 'device_id') */SKEW Hints and AQE Splitting
two_stage_agg.pyfrom pyspark.sql import functions as FTwo-Stage Aggregation and Hashed Repartitioning

Key takeaways

1
Everest histogram (one task 10x+ the median) is skew's mugshot
confirm with duration and byte ratios, never with memory bumps.
2
Census the stage's exact input
NULLs and placeholders are the usual celebrities, and the top share picks the weapon.
3
Salt hot joins with N from hot-bytes/200MB
suffix the big side, replicate the small side, restore grain carefully.
4
SKEW hints plus AQE splitting automate moderate skew; hand-salt only what automation provably can't split.
5
No-join skew needs two-stage (key, salt) partials carrying sum and count
never average partial averages.
6
Gate skew in CI (5% key-share fails the merge) and alert on task ratios in prod
drift guarantees recurrence.

Common mistakes to avoid

5 patterns
×

Enabling speculation on a skewed stage

Symptom
Duplicate copies read the same giant partition and die identically — double I/O, double cost, zero information, same 99% stall
Fix
Check the longest-to-median ratio first. Above ~10x, speculation is banned — salt or lane the key instead.
×

Raising shuffle partitions to 'fix' extreme single-key skew

Symptom
2,000 partitions still hash one key to one partition — the giant shrinks only by dilution and stays immortal far past sane counts
Fix
Partition counts divide uniform data, never a single key. Salt (or lane) celebrities; use counts for the uniform remainder.
×

Averaging partial averages in two-stage aggregation

Symptom
Revenue averages silently wrong for a quarter — partial avg-of-avgs weights small salts equally with huge ones
Fix
Carry sum and count through stage one, divide once in stage two. Assert salted vs clean-sample grains match before shipping.
×

Salting both sides symmetrically when only one side skews

Symptom
N-squared shuffle explosion: both sides replicated for no reason, job slower than the skew it treated
Fix
Salt the skewed side only; replicate or broadcast the small side. Census each side separately before choosing.
×

Forgetting NULLs in the key census

Symptom
Distinct-count looks healthy while NULL — one hash bucket holding 13% of rows — hides in the aggregation you didn't run
Fix
Count NULL share explicitly on every join key. Hash partitioners route all NULLs identically: they're the commonest celebrity.
INTERVIEW PREP · PRACTICE MODE

Interview Questions on This Topic

Q01JUNIOR
A job sits at 199/200 tasks for 3 hours. What's your first hypothesis an...
Q02SENIOR
When do you salt versus broadcast versus separate-lane a hot key?
Q03SENIOR
How does AQE skew splitting work, and where does it give up?
Q04SENIOR
Why is averaging partial averages wrong, and what's the correct pattern?
Q05SENIOR
Design a CI gate that makes skew unshippable. What does it check?
Q01 of 05JUNIOR

A job sits at 199/200 tasks for 3 hours. What's your first hypothesis and first query?

ANSWER
Skew. First the task table's longest-to-median duration and byte ratios, then a GROUP BY key census on the stage input to name the celebrity key and its share.
FAQ · 6 QUESTIONS

Frequently Asked Questions

01
How do I know it's skew and not a slow node?
02
Why does one NULL key break hash partitioning?
03
What N should I salt with?
04
Can AQE replace hand-salting entirely?
05
Does repartitioning fix skew by itself?
06
How do skewed aggregations differ from skewed joins?
N
Naren Founder & Principal Engineer

20+ years shipping production backend systems. Everything here is grounded in real deployments.

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 OutOfMemoryError in Executor During Shuffle
3 / 5 · Spark
Next
Spark AnalysisException: Cannot Resolve Column Name
→