A Spark stage finishes when its slowest task finishes. If 199 tasks take twenty seconds and one takes twenty-five minutes, the stage takes twenty-five minutes, the cluster idles for most of it, and the one executor running the big task may spill or run out of memory. That is data skew: work divided by key, where a few keys own most of the data.
This article treats skew as a problem to diagnose and fix. It covers where skew comes from, how to prove which keys are hot, what Adaptive Query Execution (AQE) does automatically and where it stops, and the manual techniques for everything else. AQE as a whole is covered in Spark AQE architecture and join strategy selection in Spark join hints; this page stays with skew. Config names and defaults below are taken from the Spark 4.2.0 performance tuning documentation.
Why one key can stall a whole stage
Every wide operation (join, groupBy, window, repartition by column) shuffles rows so that all rows with the same key land in the same partition. The partition is chosen by hashing the key modulo the number of shuffle partitions. Hashing spreads distinct keys evenly, but it cannot split a single key: every row for customer 42 goes to one partition, and one task processes it.
Real data is rarely uniform. A marketplace has a house account that appears on millions of orders, a log table has a default user id for anonymous traffic, an event table has a null where the join key was never filled in. The partition holding that key is many times larger than the rest, so its task reads, sorts and holds far more and is the one that spills. Skew costs you twice: wall-clock time, and failures such as out-of-memory errors that only ever hit one task.
Adding executors does not help. The hot partition is still one task on one core. The fix is always to change how the work is divided.
Where skew shows up
| Operation | What is skewed | Typical cause |
|---|---|---|
| Join | Shuffle partitions of one or both sides | A hot foreign key, or nulls and sentinel values such as 0 or 'unknown' |
| groupBy / aggregation | Rows per grouping key | Popular pages, large tenants, bot traffic |
| Window functions | partitionBy groups | One user or device with millions of events |
Writes with partitionBy | Rows per output partition value | Repartitioning by a date column before writing, so each date is one task |
| Input scan | Input splits | Large unsplittable files such as gzip CSV; one task per file |
Input skew happens before any shuffle and has nothing to do with keys; fix it with splittable formats such as Parquet or ORC, or a repartition after the read. The rest of this article is about key skew.
Diagnosing: prove it before you fix it
Start in the Spark UI. Open the slow stage and read the summary metrics table, which gives min, 25th percentile, median, 75th percentile and max for task duration and shuffle read size. Skew looks like a max that is tens or hundreds of times the median while the 75th percentile stays close to the median. If all the percentiles are high, the problem is too few partitions or too much data, not skew. A walkthrough of these pages is in Spark UI.
Then find the keys. Count rows per key and sort; on very large tables sample first, since hot keys survive sampling by definition.
from pyspark.sql import functions as F
# exact, one extra shuffle of just the key column
(orders.groupBy("customer_id").count()
.orderBy(F.desc("count"))
.show(20, truncate=False))
# cheaper first look on a huge table: 1% sample, scale counts by 100
(orders.sample(fraction=0.01, seed=1)
.groupBy("customer_id").count()
.orderBy(F.desc("count"))
.show(20, truncate=False))
# do not forget nulls; they hash to one partition
orders.filter(F.col("customer_id").isNull()).count()Write down the top keys and their share of rows; across 200 partitions, a key holding 2 percent is already four times an average partition. Finally, check the plan after execution: when AQE handled a skewed join, the final adaptive plan marks the sort-merge join as skewed and the SQL tab shows skewed partition counts on the shuffle read node. If the marker is absent, AQE did not act and you need to find out why.
What AQE does automatically
With AQE on (the default since Spark 3.2), Spark re-plans at every shuffle boundary using real map output sizes. One rule, OptimizeSkewedJoin, looks at the partitions feeding a sort-merge join. A partition counts as skewed only if it is larger than skewedPartitionFactor times the median partition size and larger than skewedPartitionThresholdInBytes. A skewed partition is split into smaller pieces, and the matching partition on the other side is read again for each piece, so every piece still sees all the rows it can join with.
| Config | Default | Meaning |
|---|---|---|
spark.sql.adaptive.skewJoin.enabled | true | Split skewed partitions in sort-merge joins |
spark.sql.adaptive.skewJoin.skewedPartitionFactor | 5.0 | Multiple of the median a partition must exceed |
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes | 256MB | Absolute size a partition must also exceed |
spark.sql.adaptive.advisoryPartitionSizeInBytes | 64MB | Target size for the split pieces (and for coalescing) |
spark.sql.adaptive.forceOptimizeSkewedJoin | false | Apply the rule even if it needs an extra shuffle (since 3.3) |
Three details catch people. First, both thresholds must be crossed, so whichever is larger binds: with a 40 MB median, 5x is 200 MB and the 256 MB floor binds; with a 120 MB median, 5x is 600 MB and the factor binds. Lower the factor or the threshold if you see a straggler AQE ignored. Second, coalescePartitions.parallelismFirst defaults to true, which makes coalescing ignore the 64 MB advisory size; the advisory size still sets the target for skew splits. Third, splitting happens at map-output boundaries, so a partition can be cut into at most as many pieces as there were upstream map tasks.
What AQE does not cover: the documentation describes skew handling for sort-merge joins. Broadcast joins have no shuffle and so no skew to split. Aggregations, window functions and skewed writes are not split by this rule. Some join shapes cannot be split on the skewed side without breaking semantics; for example, the preserved side of an outer join can be split, but splitting the other side would produce duplicate unmatched rows. If AQE leaves a straggler in place, use the manual techniques below.
Salting a skewed join
Salting gives a hot key many identities. Add a small random or hashed integer, the salt, to the key on the large side, and replicate the small side once per salt value. The join on (key, salt) is now spread across as many partitions as there are salt values. Salt only the hot keys, so the small side grows by a few rows instead of multiplying by the salt count.
from pyspark.sql import functions as F
SALT = 32
hot = [r.customer_id for r in (orders.groupBy("customer_id").count()
.orderBy(F.desc("count")).limit(10).collect())]
is_hot = F.col("customer_id").isin(hot)
# large side: deterministic salt for hot keys, 0 for everyone else
orders_s = orders.withColumn(
"salt", F.when(is_hot, F.pmod(F.hash("order_id"), F.lit(SALT))).otherwise(F.lit(0)))
# small side: hot keys exploded to SALT copies, others keep salt 0
customers_s = customers.withColumn(
"salt", F.explode(F.when(is_hot, F.array(*[F.lit(i) for i in range(SALT)]))
.otherwise(F.array(F.lit(0)))))
joined = orders_s.join(customers_s, ["customer_id", "salt"]).drop("salt")Use a hash of an existing unique column for the salt rather than rand(). If a task is retried and its input arrives in a different order, a random salt can send a row to a different partition the second time; Spark tries to detect such indeterminate stages and rerun them, but a deterministic salt removes the question. Pick the salt count so the hot key's share divided by the salt count is close to an ordinary partition.
Two-stage aggregation for skewed groupBy
For sums, counts, min and max, Spark already does a partial aggregation on the map side before the shuffle, so a hot key arrives as one partial row per map task. That is why skewed groupBy is often milder than skewed join. It still hurts when the map side cannot reduce: collect_list, exact distinct counts, percentiles, or keys whose rows are all distinct. Then aggregate twice: first by (key, salt), then by key.
SALT = 64
partial = (events
.withColumn("salt", F.pmod(F.hash("event_id"), F.lit(SALT)))
.groupBy("page_id", "salt")
.agg(F.count("*").alias("n"),
F.approx_count_distinct("user_id").alias("u_approx"))) # see note
final = partial.groupBy("page_id").agg(F.sum("n").alias("events"))The second stage must combine the partials correctly. Counts and sums add; averages need sum and count carried separately; distinct counts do not add (the same user can appear under several salts), so either salt by the distinct column itself, which makes the partial sets disjoint, or switch to a mergeable sketch. The code above notes the trap rather than summing the approximations.
Isolate, null keys and skewed writes
Split and union. When a handful of keys dominate, process them separately. Filter the hot keys out of both sides, join the rest normally, join the hot rows with a broadcast of the (small) matching dimension rows, and union the two results. Each half gets the right strategy, and the plan is easy to read. The same idea in Hive is covered in Hive skew join optimization.
Null and sentinel keys. In an inner join, null keys never match, so filter them out before the shuffle. In an outer join they must be kept: replace the null with a value that cannot match, spread over many partitions, for example a negative number derived from a hash of the row id, and restore null afterwards. Sentinels such as 0 or 'unknown' deserve the same treatment if they are not real matches.
Skewed writes. df.repartition("dt").write.partitionBy("dt") sends each date to a single task, so a busy day becomes one giant task and one giant file. Repartition by the date plus a bounded salt, or use the REBALANCE hint, which lets AQE split oversized output partitions (rebalancePartitionsSmallPartitionFactor controls merging of small leftovers). Cap file size with the maxRecordsPerFile write option.
(events
.repartition("dt", F.pmod(F.hash("event_id"), F.lit(8))) # up to 8 tasks per date
.write.option("maxRecordsPerFile", 5_000_000)
.partitionBy("dt").mode("append").parquet(path))Window functions. No automatic split exists for Window.partitionBy. Reduce the data first (filter bots, pre-aggregate), or restructure: compute per-(key, day) results and combine them in a second, much smaller window.
Worked example
A nightly job joins about 33 GB of orders with customers on customer_id using 200 shuffle partitions. The median partition is about 120 MB and a typical task takes 20 seconds, but the stage takes 25 minutes. The UI shows a max shuffle read of 9 GB, and the key count finds one marketplace house account owning about 27 percent of orders.
Would AQE split it? The factor rule needs more than 5 x 120 MB = 600 MB, and the threshold needs more than 256 MB; 9 GB passes both. Suppose, though, that the final plan shows the sort-merge join with no skew marker. With both thresholds cleared, the remaining documented reason is that applying the rule would introduce an extra shuffle, which AQE avoids by default; forceOptimizeSkewedJoin exists for exactly that case. Two choices: set it and accept the extra shuffle, or salt. Forcing it splits the 9 GB partition into roughly 140 pieces near 64 MB (bounded by the number of map tasks), each taking about 11 seconds, and the stage drops to the time of the normal tasks plus scheduling, around a minute in this example.
An explicit split-and-union is the alternative: route the house account's rows to a broadcast join against its single customer row. The fix is then visible in code and independent of AQE heuristics.
Failure modes
- Tuning partitions instead of keys. Raising
spark.sql.shuffle.partitionsmakes normal partitions smaller but leaves the hot key in one task. - Assuming AQE fired. It only acts on sort-merge joins whose partitions cross both thresholds; read the final plan.
- Salting both sides randomly. Rows only match when the salt matches; the small side must be replicated across every salt value.
- Salt explosion. Salting every key multiplies the small side by the salt count; salt only the hot keys.
- Wrong second-stage merge. Summing distinct counts or averaging averages produces plausible but wrong numbers.
- Hot keys move. A list of hot keys computed once goes stale; recompute it each run or derive it from the data in the same job.
Trade-offs
| Technique | Cost | Best when |
|---|---|---|
| AQE skew join | None to configure; extra reads of the other side | Sort-merge joins with a few oversized partitions |
| Broadcast the small side | Driver and executor memory | One side fits comfortably in memory |
| Salting | Code complexity, replicated rows | Joins AQE will not split, very hot keys |
| Two-stage aggregation | An extra shuffle | Aggregations that cannot reduce on the map side |
| Split and union | Two plans to maintain | A short, stable list of hot keys |
For a deeper look at what each shuffle costs, see Spark shuffle.
What to do next
- Open your slowest stage in the Spark UI and compare max to median for duration and shuffle read size.
- Run the top-20 key count on the join or grouping key, including a separate null count.
- Confirm AQE is enabled and read the final adaptive plan for a skew marker on each sort-merge join.
- If a straggler remains, check the thresholds against your actual median partition size and lower them if needed.
- For joins AQE will not split, salt only the hot keys with a deterministic salt, or split and union with a broadcast.
- For skewed writes, repartition by the output column plus a bounded salt and set maxRecordsPerFile.
- Add a check to the job that logs the top keys and their share each run, so you notice when skew changes.