Apache Spark — The answer to a slow job is in the plan and the event log
When one task never finishes, the cause is one key
In one line
A shuffle sends the same key to the same task, so if one key takes a large share of the data, the single task that received that key holds up the whole stage. This is skew, and there are three ways to solve it: AQE skew join, salting and splitting out hot keys.
Why one task is a problem
A stage finishes only when the slowest task finishes. Even if 199 tasks finish in 2 seconds, if one takes 3 minutes, the stage takes 3 minutes. Meanwhile the rest of the cores sit idle. Doubling the executors does not help: the one slow task still runs on one core.
Why is only one slow? The RDD programming guide describes a shuffle as an all-to-all operation that reads all the partitions in order to gather all the values of a key in one place. Whether it is a join or an aggregation, the destination partition is decided by the hash of the key, so the same key always goes to the same task. This is a rule that guarantees the correctness of the shuffle and, at the same time, the cause of skew. If one bot made 96,000 of 240,000 clicks, 40% piles into the one partition that received that user ID. Increasing the number of partitions does not solve it, because a single key cannot be split.
How to recognize it
Skew is not visible in averages. The bytes read by the whole stage look fine. What you must look at is the per-task distribution. The stage detail screen described in the web UI documentation has summary metrics for all the tasks, and among them Duration and Shuffle Read Size / Records are shown as minimum, median and maximum. How many times the median the maximum is is the size of the skew. If it is tens of times, it means one task is doing almost all the work alone.
The same numbers are in the event log too. The record left each time a task ends contains the number of records that task read through the shuffle, so with just one log file you can count the maximum and median directly without the UI.
The first way: AQE skew join
The skew join optimization in the performance tuning documentation splits the skewed partitions of a sort merge join into several tasks of similar size, and replicates the matching partition on the other side as many times as needed. Both spark.sql.adaptive.enabled and spark.sql.adaptive.skewJoin.enabled must be on, and both default to true.
Which partitions are skewed is decided when both conditions are met. The size must be greater than skewedPartitionFactor times the median (default 5.0) and, at the same time, greater than skewedPartitionThresholdInBytes (default 256MB). Because of the second condition, on small lab data, no matter how skewed it is, nothing happens with the default settings. That is why the lab lowers the threshold. The size targeted when splitting is advisoryPartitionSizeInBytes (default 64MB).
When it works, SortMergeJoin(skew=true) appears in the final plan, and the shuffle read beneath it changes to AQEShuffleRead … coalesced and skewed. There is one caveat. If the split would need to create one more shuffle, AQE by default does not apply it. What you turn on when you want to do it anyway is spark.sql.adaptive.forceOptimizeSkewedJoin (default false).
The second way: salting
Skew where there is no AQE, or outside of joins, is solved by people. Salting attaches a random number (the salt) between 0 and N-1 to the key on the large side, making one hot key into N different keys. Then the hash scatters them over N partitions. In exchange, the small side has to be replicated for every salt value so that no match is lost.
from pyspark.sql import functions as F
N = 8
clicks_s = clicks.withColumn("salt", (F.rand(7) * N).cast("int"))
users_s = users.crossJoin(spark.range(N).withColumnRenamed("id", "salt"))
joined = clicks_s.join(users_s, ["user_id", "salt"]).drop("salt")
The cost is clear: the small side becomes N times larger. So set N only as large as needed to resolve the skew, and use this when the small side is really small.
The third way: splitting out hot keys
If the hot keys are fixed to a few, there is a simpler way. You filter out only the rows of those keys and process them separately, join the rest as usual, and then attach the two with a union. On the hot-key side, only a few rows of the other table correspond to that key, so if you join with a broadcast there is no shuffle at all. The rest is evenly distributed data from which the skew has disappeared. Finding the hot keys is enough with counting per key and looking at the top few.
Which of the three to choose
The order is usually this. First, check in the final plan whether AQE is already solving it. If it is a sort merge join and the skewed partition exceeds the two conditions, it is solved without changing a single setting. If you imagine the way AQE splits, you can also see its limits. It divides the skewed side's partition into several pieces and attaches the whole matching partition from the other side to each piece. If the matching partition on the other side is also large, the replication cost grows, and if both sides are skewed on the same key, splitting the pieces does not reduce the number of matches itself.
If AQE does not solve it, look at how many hot keys are fixed. If you can name them, like one bot account or three large customers, splitting them out is the simplest and the result is also easy to explain. If the hot keys change every day or number in the dozens, salting is better, because it spreads all keys evenly even if you do not know which keys are hot. Whichever way, when you finish, measure the per-task maximum and median again to confirm that the ratio really went down.
When skew is hidden: partial aggregation
If you run groupBy("user_id").count() on the same bot data, strangely skew is barely visible. Looking at the plan, there is a reason. There is a HashAggregate(partial_count) before the shuffle, so each map task sends only one line of partial total per key. The bot's 96,000 rows have been reduced, before the shuffle, to a few numbers as many as the map tasks. This is also the reason the RDD guide recommends reduceByKey and aggregateByKey instead of groupByKey for sums and averages per key.
So skew shows up in joins and in operations without partial aggregation. A window function has to gather all the rows of a partition key into one task, and an aggregation that collects a list has partial results as large as the original rows themselves. You must not assume that because the aggregation was fine, the join will be fine too.
What it looks like in the field
First, the progress bar stops at 99%. Only one task remains and runs for tens of minutes. If that task's shuffle read records are tens of times those of the others, it is skew.
Second, the null key is a hot key. If millions of nulls pile up in a column with missing values, null is still one key by hash, so it gathers in one partition. In a join condition, null equals no value, so no match arises. But an outer join has to keep the rows without a match in the result, so those rows are shuffled as they are. If you set the null-key rows aside first and attach them with a union after the join, the result is the same and the skew disappears.
Third, skew grows. From the day a bot or a large customer appears, a job that was fine until yesterday slows down. The code is unchanged, so without looking at the data distribution you cannot find the cause. If you record the row counts of the top few join keys every day, you can spot signs of a growing hot key before the job slows down.
What really matters in practice
- Look at the per-task maximum and median, not the average. That ratio is the size of the skew.
- Increasing the number of partitions does not split a single key.
- AQE skew join works only on sort merge joins, and only when both conditions are exceeded.
- Salting pays the price of inflating the small side N times. Keep N only as large as needed.
- An aggregation with partial aggregation hides skew. Check again in joins and window functions.
What you will do in the next lab
In data where one bot made 40% of the clicks, you first find the hot user. You turn off broadcast and AQE, join with a sort merge, then pull the per-task record counts of the tasks that read the shuffle from the event log and compute the maximum and the median. Next you turn on the AQE skew join, lowering the skew threshold to fit the small lab data, and confirm that the skew marker appears in the final plan. You solve the same join again with salting and with splitting out the hot keys and compare whether the answers of the three remedies are the same, and finally confirm with the same two numbers that a per-user groupBy aggregation has almost no skew thanks to partial aggregation.