Data skew and salting
Why one task runs for an hour while the rest finish in seconds, and how salting spreads a hot key.
On this page
Show code in
Every code block on the page follows this.
You will learn
- How to recognise skew in the Spark UI
- Why one hot key makes one task run for hours
- How salting spreads a hot key across many tasks, for joins and aggregations
- Cheaper fixes to try before salting
Read first
- Jobs, stages and tasks · 3 min read
- Narrow vs wide transformations · 3 min read
Comfortable with these? Read on.
The symptom
A stage shows 199 of 200 tasks finished in seconds and one still running after an hour. In the stage detail, the max task duration and shuffle read size are orders of magnitude above the median. That one task may also spill heavily or die with an out-of-memory error.
| Task metric | Median | Max |
|---|---|---|
| Duration | 6 s | 58 min |
| Shuffle read | 120 MB | 48 GB |
| Spill (disk) | 0 | 31 GB |
The cause
Hash partitioning sends every row of a key to the same partition: that is what makes joins and aggregations correct. Real data is rarely uniform. Common culprits: null or empty keys (often millions of rows with customer_id = null), default values like "unknown" or 0, and genuinely popular entities, such as one huge customer or a viral product.
df.groupBy("key").count().orderBy(F.desc("count")).show(20). Often the top key is null, and the right fix is to filter or handle nulls separately rather than salt.Try these first
- Handle null keys separately. Null keys never match in a join, so filter them out before the join and union them back afterwards if needed.
- Broadcast the other side. If the table being joined is small enough, a broadcast join has no shuffle, so no skew.
- Let AQE split it. For sort-merge joins, AQE skew handling splits oversized partitions automatically (see AQE skew-join handling).
- Pre-aggregate. If the hot key only needs a total, aggregate before joining so the hot key becomes one row.
Salting a join
Step 1 · salt the big side
The command
N = 16 big_s = big.withColumn("salt", (F.rand() * N).cast("int"))
SELECT *, CAST(rand() * 16 AS INT) AS salt FROM big
Step 2 · replicate the small side
The command
small_s = small.withColumn("salt", F.explode(F.sequence(F.lit(0), F.lit(N - 1))))
SELECT s.*, salt FROM small s LATERAL VIEW explode(sequence(0, 15)) t AS salt
Step 3 · join on key + salt
The command
joined = big_s.join(small_s, ["key", "salt"]).drop("salt")
SELECT b.*, s.value FROM big_salted b JOIN small_salted s ON b.key = s.key AND b.salt = s.salt
Cost · the trade-off
Salting an aggregation
AQE does not split skewed aggregations, so salting is the main tool here. Aggregate in two phases: first by (key, salt), which spreads the hot key across N tasks, then by key alone over the much smaller partial results.
partial = (events
.withColumn("salt", (F.rand() * 16).cast("int"))
.groupBy("key", "salt").agg(F.sum("amount").alias("s"), F.count("*").alias("n")))
result = partial.groupBy("key").agg(F.sum("s").alias("total"), F.sum("n").alias("rows"))
SELECT key, SUM(s) AS total, SUM(n) AS rows FROM ( SELECT key, salt, SUM(amount) AS s, COUNT(*) AS n FROM (SELECT *, CAST(rand() * 16 AS INT) AS salt FROM events) GROUP BY key, salt) GROUP BY key
This works for any aggregate that can be combined: sums, counts, min, max, and averages computed as sum divided by count. It does not work directly for exact distinct counts or medians.
collect_list, and for joins.Common mistakes
Salting without finding the hot keys
Salting both sides randomly
Picking a huge N
Key takeaways
- Skew means one key, and so one task, gets a disproportionate share of rows.
- The Spark UI shows it as a huge gap between median and max task metrics.
- Handle null keys, broadcast or rely on AQE before salting.
- Salting a join: random salt on the big side, replicate the small side; for aggregations, aggregate twice.
Check yourself
3 questions1. In a salted join, what happens to the smaller table?
Show the answer
It is replicated once per salt value. Each row must meet every possible salt from the big side, so it is copied N times.
2. Which skew can AQE NOT split automatically?
Show the answer
A skewed groupBy aggregation. AQE skew handling targets sort-merge joins; aggregations need salting or other fixes.
3. The hottest join key is null. What is the simplest fix?
Show the answer
Filter null keys before the join and handle them separately. Null keys never match in an equality join, so they can be removed before shuffling.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Go deeper
AQE skew-join handlingSpark internals
Broadcast joinsSpark internals
Jobs, stages and tasks
Primary sources: Performance tuning