Salting spreads one hot key across many partitions by adding a random number to the key. You do it in three steps: salt the big side, replicate the small side, then join on the key and the salt.
Example: one customer owns most of the rows
events has 1 billion rows, and 300 million have customer_id = 7. customers is smaller, with one row per customer. A normal join sends all 300 million rows for customer 7 to a single task.
Pick the salt count N, say 10.
from pyspark.sql import functions as F
N = 10
# 1. big side: random salt 0..N-1
events_s = events.withColumn("salt", (F.rand() * N).cast("int"))
# 2. small side: one copy of each row per salt value
salts = spark.range(N).withColumnRenamed("id", "salt")
customers_s = customers.crossJoin(salts)
# 3. join on key + salt
joined = events_s.join(customers_s, ["customer_id", "salt"])Now customer 7's 300 million rows are split into 10 groups of about 30 million, and each group lands on a different partition. Each salt value on the customers side has a matching copy, so no rows are lost.
Choosing N
Compare the hot key's size with a normal partition. If the hot key is 40 times the median, N around 40 would even it out. Too small leaves skew. Too large multiplies the small side more than needed.
The cost
The small side now has N times as many rows. For 10 million customers and N = 10 that is 100 million rows. To reduce this, salt only the hot keys: give hot keys a random salt, and give all other keys salt 0. Replicate only the hot keys' small-side rows N times.
Salting an aggregation
For a group-by on a skewed key, use two stages: aggregate by (key, salt) first, then aggregate again by key alone.
partial = df.groupBy("customer_id", "salt").agg(F.sum("amount").alias("s"))
final = partial.groupBy("customer_id").agg(F.sum("s").alias("total"))This works for sums and counts. For averages, carry sum and count and divide at the end. For distinct counts it needs more care.
Try these first
AQE skew join handling often solves it without code changes. Salting is for when AQE is not enough, or on aggregations and operators it does not cover.