TT Lab
Get started
Learn Learning paths Courses

Apache Spark — The answer to a slow job is in the plan and the event log

Measure and fix skew in clicks where one bot makes 40%

Continue in TT Lab

Goal

You join data in which one bot made 40% of the clicks with the users table, and measure with the event log how records pile up in a single post-shuffle task. You solve it in three ways, AQE skew join, salting and splitting out the hot key, and also confirm why the same hot key is not a problem in a count aggregation.

Why it matters

Most reports that "the job stops at 99%" are skew. A shuffle decides partitions by the hash of the key, so the same key always goes to one task. If a single key is 40% of the whole, one task takes on 40%, and even after all the other tasks have finished, only that one keeps running. Adding more cores does not help: that task cannot be split. Skew is more accurately judged by distribution than by time. If you compare the maximum and the median of the records read per task inside a stage, you get the same numbers whether the machine is fast or slow. This lab also judges by that ratio, not by time. There are three remedies. AQE looks at the actual size after the shuffle, splits a large partition into several pieces and replicates the matching partition on the other side (if the settings are right, there is no code change). Salting attaches a random tag to the hot key to scatter it over several partitions and inflates the other side by that number. The simplest is to split out the hot key and process it separately. And there are aggregations where skew is not a problem to begin with, because there is partial aggregation that reduces on the map side first.

Steps

  1. Put a function that reads the clicks and users tables in /root/spk/skew/common.py, and with /root/spk/skew/hot.py (application spk-skew-hot), write the top 5 users by number of clicks (ties broken by user_id ascending) as a CSV (user_id,clicks) to /root/spk/skew/out/hot.
  2. With /root/spk/skew/plain.py (application spk-skew-plain, threshold -1, AQE off, 8 shuffle partitions), join the clicks and users and write the number of clicks and the sum of ms per segment as a CSV (segment,clicks,ms) to /root/spk/skew/out/by_segment.
  3. From the log of the step 2 application, among the stages that read the shuffle, pick the stage whose per-task maximum of records read is the largest, and write its stage number, maximum and median (the lower one) to /root/spk/skew/out/skew.json.
  4. With /root/spk/skew/aqe.py (application spk-skew-aqe, threshold -1, AQE on, 16 shuffle partitions, skew threshold 256k, advisory size 64k), write the same result to /root/spk/skew/out/by_segment_aqe.
  5. With /root/spk/skew/salt.py (application spk-skew-salt, the same settings as step 2), attach a salt of 0–7 only to the bot's clicks, inflate the bot row on the user side to eight copies, join on user_id and salt, and write the result to /root/spk/skew/out/by_segment_salt.
  6. With /root/spk/skew/split.py (application spk-skew-split, the same settings as step 2), do a broadcast join for the bot's clicks and an ordinary join for the rest, attach them with unionByName, and write the result to /root/spk/skew/out/by_segment_split.
  7. With /root/spk/skew/agg.py (application spk-skew-agg, AQE off, 8 shuffle partitions), write the number of clicks per user as Parquet to /root/spk/skew/out/per_user, and write the per-task record maximum and median (the lower one) of the stage in that application that read the shuffle to /root/spk/skew/out/agg_skew.json.
  8. In /root/spk/skew/report.md, write three sections: ## 쏠림의 모양, ## 세 가지 처방 and ## 부분 집계 (use exactly these Korean headings in this order; they mean "The shape of the skew", "Three remedies" and "Partial aggregation"). Put the two numbers from step 3 in the first section and the two numbers from step 7 in the third section.

Notes

Find the hot key

In /root/spk/skew/common.py, put a function that reads the clicks (/data/clicks/clicks.jsonl) and users (/data/clicks/users.csv) tables, create /root/spk/skew/hot.py with the application name spk-skew-hot, and write the top 5 users by number of clicks (clicks descending, ties by user_id ascending) as a CSV with a header row (user_id,clicks) to /root/spk/skew/out/hot.

The first step in solving skew is knowing which key is hot. Look at the gap between 1st and 2nd place. If 1st place is several times the rest, that single key holds up one task.

A plain join piles up in one task

Create /root/spk/skew/plain.py with the application name spk-skew-plain and the settings spark.sql.autoBroadcastJoinThreshold=-1, spark.sql.adaptive.enabled=false and spark.sql.shuffle.partitions=8, join the clicks and users by user_id, and write the number of clicks per segment (clicks) and the sum of ms (ms) as a CSV with a header row (segment,clicks,ms) to /root/spk/skew/out/by_segment.

Broadcast and AQE are turned off in order to expose the skew on purpose. A sort merge join shuffles both sides by user_id, so all of the bot's roughly 90,000 click rows go to one partition. The grader checks how many times the median the maximum is in that stage.

Skew in numbers: maximum and median

From the most recent log of the step 2 application (spk-skew-plain), collect the per-task Total Records Read for each stage that read the shuffle, pick the stage with the largest maximum, and write to /root/spk/skew/out/skew.json as {"stage_id": 정수, "max_records": 정수, "median_records": 정수} (all three values are integers). The median is the value at index (n-1)//2 (from 0) among n values in ascending order.

The join stage reads both shuffles together, so the task of the partition containing the bot reads all of the bot's clicks and one row of the bot user. Maximum ÷ median is the size of the skew. When this ratio exceeds 3, it is commonly called "skewed".

Remedy 1: AQE skew join

Create /root/spk/skew/aqe.py with the application name spk-skew-aqe and the settings spark.sql.autoBroadcastJoinThreshold=-1, spark.sql.shuffle.partitions=16, spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=256k and spark.sql.adaptive.advisoryPartitionSizeInBytes=64k (leaving AQE on), and write the same result as in step 2 to /root/spk/skew/out/by_segment_aqe.

AQE considers a partition skewed when it is several times the median (skewedPartitionFactor, default 5) and also larger than the threshold bytes. The default threshold (256MB) does not suit this small data, so it was lowered. If SortMergeJoin(skew=true) appears in the final plan, it was split without any code change.

Remedy 2: sprinkle salt on the hot key

Create /root/spk/skew/salt.py with the application name spk-skew-salt and the same settings as step 2 (threshold -1, AQE off, 8 shuffle partitions). On the clicks side, attach a random salt between 0 and 7 if the user is the bot and 0 otherwise; on the users side, inflate only the bot row to eight copies with salt 0–7 (the others get 0); then join on user_id and salt and write the result to /root/spk/skew/out/by_segment_salt.

Because the salt goes into the join key, the bot's clicks are scattered over eight partitions. The bot row on the other side must be eight copies so that there is a match whichever salt it goes with. The result must be exactly the same as in step 2, and the grader checks whether the skew ratio has dropped below half of step 2's.

Remedy 3: split out the hot key

Create /root/spk/skew/split.py with the application name spk-skew-split and the same settings as step 2, do a broadcast join of the bot's clicks with the one bot user row, do an ordinary join of the remaining clicks, attach them with unionByName, and write to /root/spk/skew/out/by_segment_split.

If there is one hot key and you know its name, this is the simplest and most reliable remedy. The other side of that key is one row, so the broadcast is free, and the rest is spread evenly. BroadcastHashJoin, SortMergeJoin and Union should all be visible in the plan.

Why count is not skewed: partial aggregation

Create /root/spk/skew/agg.py with the application name spk-skew-agg and the settings spark.sql.adaptive.enabled=false and spark.sql.shuffle.partitions=8, write the number of clicks per user (user_id,clicks) as Parquet to /root/spk/skew/out/per_user, and write the per-task record maximum and median (the lower one) of the stage in that application that read the shuffle to /root/spk/skew/out/agg_skew.json as {"max_records": 정수, "median_records": 정수} (both values are integers).

The same bot is there, but this time the tasks are divided evenly. The first HashAggregate in the plan is partial_count: each map task counts per user first, so what passes to the shuffle is one row per user per map task. Skew is a problem for operations that cannot be reduced first like this (a join, or something like collect_list).

Record the skew and the remedies in numbers

In /root/spk/skew/report.md, write three sections: ## 쏠림의 모양, ## 세 가지 처방 and ## 부분 집계 (use exactly these Korean headings in this order; they mean "The shape of the skew", "Three remedies" and "Partial aggregation"). Put the maximum and median from step 3 in the first section and the maximum and median from step 7 in the third section, as numbers.

In the second section, write one line each on what the three remedies change (code or settings, and what they replicate). It would also be good to write which you would try first in production.