Skip to content
Great engineers know 5 min read 1 practice problem ↓

Data skew and salting

Why one task runs for an hour while the rest finish in seconds, and how salting spreads a hot key.

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

Comfortable with these? Read on.

TL;DR A shuffle sends all rows of a key to one partition. If one key has 40% of the rows, one task does 40% of the work while the others idle. Salting appends a random number to the hot key so its rows spread over several partitions, then undoes the split afterwards.

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 metricMedianMax
Duration6 s58 min
Shuffle read120 MB48 GB
Spill (disk)031 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.

Tip: find the hot keys first: 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

  1. 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.
  2. Broadcast the other side. If the table being joined is small enough, a broadcast join has no shuffle, so no skew.
  3. Let AQE split it. For sort-merge joins, AQE skew handling splits oversized partitions automatically (see AQE skew-join handling).
  4. 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

Add a random salt from 0 to N-1 to every row of the large, skewed table. Rows of the hot key now carry N different (key, salt) combinations.

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

Every row of the other table must be able to meet every salt, so copy it N times, once per salt value.

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

Joining on both columns hashes the hot key into up to N partitions instead of one. Each task gets about 1/N of the hot key's rows.

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

The small side grows N times, and the join now shuffles that bigger table. Pick N from the skew: if the hot key is 50 times the median partition, N of 32 to 64 is reasonable. Salting only the hot keys (salt = 0 for everything else, and replicate only those keys) keeps the overhead small.

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.

PySparkSpark SQL
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.

Note: plain sums and counts are already partially aggregated before the shuffle, so skew hurts them less. Salting matters most for aggregations that cannot shrink early, such as collect_list, and for joins.

Common mistakes

Salting without finding the hot keys

Often the hot key is null, and filtering it is the whole fix.

Salting both sides randomly

Random salts on both sides would almost never match. Randomise one side, replicate the other.

Picking a huge N

The replicated side grows N times; keep N proportional to the skew.

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 questions

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

Solve: Find Hot Keys and Plan Salting →

Go deeper

Primary sources: Performance tuning