In a Spark join, how does salting a hot key spread its rows across many tasks?
answer
- make the key wider so it stops collapsing
- one side gets randomness, the other gets copies
- join on the pair, then drop the extra column
- salt both sides randomly and rows vanish
- size it from the hot key's bytes
basics
~20 sYou append a random 0..N-1 salt to the join key on the big side and replicate every row of the small side N times, once per salt value. The composite key now hashes into N partitions instead of one.
solid answer
~50 sSkew comes from hashing: every row carrying the hot key goes to one partition. Salting changes the key so it no longer collapses. On the large side, add a column such as `(rand() * N).cast("int")`; on the small side, explode each row into N copies, one per salt value; then join on both the original key and the salt. The hot key's rows now spread over N partitions and N tasks. The result is identical to the unsalted join because each large-side row exists once and finds its matching small-side row in the copy carrying its salt. The costs are the N-fold blow-up of the replicated side and the extra code, which is why you usually salt only the keys you have measured as hot, and pick N from the hot key's size rather than by feel.
code
python · 11 linesfrom pyspark.sql import functions as F
N = 32
big = orders.withColumn("salt", (F.rand() * N).cast("int"))
small = customers.withColumn(
"salt", F.explode(F.array([F.lit(i) for i in range(N)])))
joined = (big.join(small, ["customer_id", "salt"], "inner")
.drop("salt"))go deeper
Know the idea: a single key value that holds most of the rows can be split by adding a small random number to it, so the work lands in several tasks instead of one.
Write the transformation correctly — random salt on one side, exploded copies on the other, join on the pair — and explain why that keeps the output identical to the unsalted join.
Show the judgment around it: measure the hot key first, try broadcasting and null handling before salting, size N from real bytes, and salt only the keys that need it.
Weigh salting's permanent complexity against fixing the distribution upstream — pre-aggregating, partitioning the source differently, or isolating a dominant tenant — and decide what the platform should own versus what each pipeline reimplements.
## The problem A shuffle join hash-partitions both sides on the join key. If one key value carries a large share of the rows, all of those rows land in one partition, one task reads them all, and the stage waits. No amount of extra parallelism helps, because the key is atomic: `hash("c-42")` has exactly one answer. Salting attacks the cause. If the key is what concentrates the data, make the key artificially wider. ## The recipe 1. Choose a salt count `N`. 2. On the **large** (skewed) side, add a salt column with a random integer in `[0, N)`. 3. On the **small** side, replicate each row `N` times, once per salt value — typically by adding a literal array of `0..N-1` and exploding it. 4. Join on `(original_key, salt)` instead of `original_key`. 5. Drop the salt column from the result. After step 4 the hot key has become `N` distinct composite keys, which hash to up to `N` different partitions, so up to `N` tasks share the work. ## Why the result is still correct This is the part interviewers probe. Each large-side row is assigned exactly one salt, so it exists once. Each small-side row exists in all `N` salted copies, so whichever salt the large-side row drew, there is a matching copy carrying it. Every pair that would have matched in the unsalted join matches exactly once in the salted one — no duplicates, no drops. The classic mistake is salting **both** sides randomly. Then a large-side row with salt 7 and its matching small-side row with salt 3 never meet, and rows silently disappear from the output. "Salt one side randomly, replicate the other exhaustively" is the invariant. ## Choosing N Size it from the data, not from taste: take the hot key's total bytes and divide by the partition size you want each task to handle, rounding up. If one key holds 20 GB and you want ~200 MB tasks, `N` is about 100. Going larger costs you: the replicated side's row count is multiplied by `N`, and shuffled accordingly, so an `N` of 1000 against a 10-million-row dimension means ten billion rows in the exchange. In practice `N` between 8 and 64 covers most cases. ## Salting only the hot keys Salting every key multiplies the small side by `N` for nothing — the well-behaved 99.9% of keys were never a problem. The refined version salts only the values you have measured as hot: the large side gets a random salt for hot keys and a constant `0` for everything else, and the small side is split into a "cold" copy with salt `0` and a "hot" copy exploded across all `N` salts, then unioned. Note that a generator such as `explode` cannot be nested inside a `when` expression in Spark, which is why the small side is built by filtering and unioning rather than by a single conditional column. The drawback is that the hot-key list is a hard-coded assumption about the data. It rots. Either compute it at runtime from a cheap count of the key, or schedule a review of it — a salted job whose hot-key list is a year out of date is skewed again, and nobody notices until it misses an SLA. ## Salting an aggregation For a skewed `groupBy` — which the adaptive skew-join rule does **not** cover — the equivalent is two-phase aggregation: group by `(key, salt)` to compute partials spread across `N` tasks, then group the partials by `key` alone to combine them. This works cleanly for associative, commutative aggregates like `sum`, `count`, `min` and `max`. It does not work for aggregates whose partials cannot be merged, such as an exact median or an exact distinct count, where you need a different approach (approximate sketches, or a dedicated pass over the hot group). ## When *not* to salt Salting is real code in your pipeline, with real correctness risk, so exhaust the cheaper options first: - **Can you broadcast the small side?** That removes the shuffle entirely and no salt is needed. - **Is the hot key a null or a sentinel?** Nulls never match in an inner equi-join; filtering or handling them separately usually costs one predicate. - **Is AQE's skew-join splitting already handling it?** For a shuffle join whose skewed partition clears the factor and the byte threshold, Spark splits it for you at runtime. What salting still buys you when those fail: it changes the distribution **before** the shuffle is written, so the map side, the shuffle write and the fetch are all balanced — whereas the adaptive rule only redistributes the read of an already-written giant partition. And it covers aggregations, which the adaptive rule does not touch at all.
- How do you salt a skewed groupBy instead of a join?Two-phase aggregation. Group by (key, salt) so the hot group's rows are spread over N tasks and produce partial aggregates, then group those partials by key alone to combine them. It works for associative aggregates such as sum, count, min and max. An exact median or exact distinct count cannot be combined this way and needs sketches or a separate pass.
- How do you choose the salt count N?From measurement: divide the hot key's total bytes by the partition size you want per task and round up. If one key holds 20 GB and you want 200 MB tasks, N is around 100. Do not go larger for safety — every extra salt multiplies the replicated side's rows and the bytes shuffled, so an oversized N trades one bottleneck for another.
- Why bother salting when AQE's skew-join splitting is on by default?The adaptive rule only redistributes the read of an already-written oversized partition, and only for shuffle joins that clear both its factor and its byte threshold. Salting changes the distribution before the shuffle is written, so the map side and the shuffle write are balanced too, and it also covers aggregations, which the rule does not handle.
- What is the maintenance risk of a salted job?The salt count and any hard-coded hot-key list encode an assumption about the data distribution on the day it was written. Keys cool down, new ones heat up, and the job quietly returns to being skewed — or pays N-fold replication for nothing. Derive the hot keys at runtime where you can, and monitor the stage's max-versus-median task duration so the decay is visible.
A single service window for everyone named Smith becomes N windows labelled Smith-0 to Smith-N. Arrivals pick a window at random, and every clerk keeps a full copy of the Smith file so any window can serve them.
saying these in an interview costs you the question
- Salts both sides randomly, silently dropping matching rows
- Salts every key rather than the measured hot ones
- Thinks salting removes the shuffle instead of spreading it
- Picks the salt count without measuring the hot key
- Claims the salted join returns different results than the original