This is data skew. One partition holds far more data than the others, so one task does most of the work while 199 cores sit idle. The stage can only finish when its slowest task does.
Confirm it first
Open the slow stage in the Spark UI and look at the task summary table. In a skewed stage, the max for shuffle read size and records is many times the median, and max duration is far above the 75th percentile. If sizes are even but one task is slow, suspect a bad node instead.
Find the hot key
(df.groupBy("customer_id").count()
.orderBy(F.desc("count")).show(10))Very often the winner is NULL, an empty string, or a default id like -1 or unknown. A million rows from guest checkouts with customer_id = NULL all hash to the same partition.
Fixes, from cheapest to most work
- Filter out junk keys if they have no meaning for the join, such as NULL keys in an inner join (they never match anyway).
- Turn on AQE skew join handling. With
spark.sql.adaptive.enabledandspark.sql.adaptive.skewJoin.enabled(both on by default in recent versions), Spark splits a partition that is much larger than the median (by default more than 5 times the median and above 256 MB) into smaller pieces at runtime. It handles sort-merge joins but not every operator, so check that it kicked in. - Broadcast the small side, if it is small enough, which removes the shuffle completely.
- Isolate the hot keys: process the few giant keys separately (for example with a broadcast join), process the rest normally, then union the results.
- Salting: add a random suffix to the hot key so it spreads over many partitions (see the salting question).
- For aggregations, do a two-stage aggregation: partial aggregate with a salt, then the final one.
Stop it coming back
Add a check on the top key share in your pipeline, or alert on max-to-median task time. Skew grows silently as one customer or one event type grows. Say that you would fix NULL keys at the source if they come from a bug, since that is the cheapest long-term fix.