A hash join's build-side hash table does not fit in the memory budget for that operator. Explain what the engine does instead, and what goes wrong when the join key is heavily skewed.
answer
- Grace hash join: partition both sides by same hash
- Equal keys -> same partition, so no cross-partition matches
- Extra cost ~ one write + one read of both inputs
- Partition too big -> repartition with a new seed
- One dominant key can't be split; sentinel values are the usual culprit
basics
~20 sIt partitions both inputs by a hash of the join key into batches written to temporary files, then joins one pair of matching partitions at a time in memory. If one key value dominates, its partition stays too big to fit, so partitioning recurses without helping and that batch degrades badly.
solid answer
~60 sThe engine falls back to a **partitioned (grace) hash join**. Using a hash function on the join key, it splits the build input into P partitions, keeping one in memory and writing the rest to temporary files; the probe side is partitioned with the *same* function so matching rows always land in the same partition number. Then it processes partition pairs one at a time: build the hash table for partition i, probe it with the probe side's partition i, emit matches, discard, move on. Correctness holds because equal keys hash equally, so no match can straddle partitions. The cost is roughly one extra write plus read of both inputs. If a partition still exceeds the budget, the engine repartitions it with a different hash seed - **recursive partitioning**. Skew breaks this. If one key value covers a large share of rows, every row with that value hashes to the same partition, so repartitioning cannot split it: recursion just burns I/O. The batch is processed with repeated probe passes or by degrading to a nested-loop-style strategy, and the join can take orders of magnitude longer while consuming large temporary space.
go deeper
Know that a hash join builds a lookup table on one input and probes with the other, and that too much data means writing parts to temporary files.
Describe grace hash join: partition both inputs with the same hash, process matching partition pairs, and pay one extra write and read.
Cover recursive partitioning, why a single dominant key defeats it, the observable symptoms, and remedies from filtering sentinels to separating the hot key.
Treat skew as a data and modelling issue as much as an execution one: sentinel design, statistics quality, distribution strategy, and when to accept a different join or plan shape.
## The in-memory case first A hash join runs in two phases: **build** a hash table over the smaller input keyed by the join column, then **probe** it by streaming the other input and looking up each row's key. If the build side fits in the operator's memory budget, the cost is roughly one read of each input - which is why hash join is the workhorse for large equality joins. Everything below is about what happens when the build side does not fit. ## Partitioned (grace) hash join The engine does not give up and switch to a different join algorithm; it makes the problem smaller. 1. Choose a partition count P so that an expected build partition fits the budget, based on estimated build size divided by budget. 2. Scan the build input, apply a hash function to the join key, and route each row to partition `h(key) mod P`. One partition is often kept in memory; the rest are written to temporary files (**batches**). 3. Scan the probe input, applying the **same** hash function, writing each row into the matching partition file. Rows belonging to the in-memory partition are joined immediately. 4. For each remaining i: read build partition i, construct the hash table, stream probe partition i through it, emit matches, free the table. Correctness rests on one property: identical keys always produce identical hashes, so a build row and a matching probe row necessarily land in the same partition number. No cross-partition comparison is ever needed. The extra cost is bounded and sequential: each input is written once and read once from temporary files, on top of the original scans. Note that it depends on hash-partition placement, not on the natural row order, so partitioning is not a sort. ## When a partition still does not fit Because partition sizes are estimated, one may still exceed the budget - the row estimate was wrong, or rows are wider than assumed. The engine then splits that partition again using a **different hash seed**, so rows with different keys that collided under the first function separate under the second. This is recursive partitioning, and it usually converges after one extra level. ## Where skew breaks it Recursive partitioning only works if the partition contains *many distinct keys*. Suppose 40% of the probe rows carry the same value - a sentinel like `-1`, a default tenant, an `UNKNOWN` dimension member, or the busiest customer. Every one of those rows hashes to one partition, under every seed. Repartitioning produces one huge partition and several empty ones, forever. What the engine does then varies, and none of it is good: - **Multiple probe passes:** load as much of the oversized build partition as fits, stream the entire probe partition against it, then repeat for the next slice. With s slices the probe partition is read s times. - **Degrade to a nested-loop-shaped strategy** for that batch. - **Explode temporary space**, since both sides of a giant partition are staged on disk. Symptoms in practice: the join's actual row counts wildly exceed estimates, the number of batches climbs during execution, temporary file usage spikes, and one operator dominates runtime while the rest of the plan looks healthy. CPU stays busy on one worker while others idle, because a parallel hash join distributes by hash and the skewed value pins one worker. ## What to do about it - **Fix the data first.** Sentinel values that mean "unknown" are the classic cause. Excluding or NULL-ing them, or filtering them before the join, often removes the entire problem, because rows that cannot match should not be partitioned at all. - **Handle the hot key separately** - join the dominant value with its own strategy and union the result with the join over the remaining keys. - **Improve statistics.** Skew is exactly what per-value or histogram statistics exist to describe; without them the optimizer sizes the join for an average key and is guaranteed to be wrong. - **Reconsider the plan shape.** With a very skewed large side, broadcasting a small build side (replicating it to every worker) avoids hash distribution entirely, and a merge join over already-ordered inputs sidesteps hash partitioning. - **Raise the operator's memory budget narrowly**, for that workload only, if the shortfall is genuinely modest. The general lesson: spilling is a graceful, linear degradation for evenly distributed data, but it is *not* graceful for skewed data, because the mechanism it relies on - splitting by hash - cannot split a single key.
- Why must the probe side be partitioned with exactly the same hash function as the build side?Because the join's correctness depends on every matching pair meeting in the same batch. Equal keys hash to the same value only if the same function and modulus are applied on both sides; with different functions a build row and its matching probe row could land in different partitions and the match would silently be lost. It is also why recursive partitioning must apply its new seed to both sides of that partition.
- How does skew on a join key show up differently in a parallel hash join than in a single-threaded one?A parallel hash join distributes rows across workers by hash of the join key, so all rows with the dominant value are routed to one worker. That worker's runtime and memory demand dwarf the others, and the query's duration becomes the straggler's duration while the remaining workers sit idle. The tell is a large gap between the slowest and average worker times, with one worker also holding most of the temporary file usage.
saying these in an interview costs you the question
- Saying the engine switches to a nested loop join whenever the hash table does not fit, missing partitioned/grace hash join
- Believing partitioning sorts the data, or that matches must be checked across partitions
- Claiming recursive repartitioning always resolves an oversized partition, ignoring single-key skew
- Treating a spilling hash join as evidence that hash join was the wrong algorithm, rather than checking estimates and skew
- Ignoring sentinel or default key values as a source of skew