skip to content

What is a grace hash join (and its hybrid variant), and how does an execution engine complete a hash join correctly when the build input does not fit in the memory budget?

level: seniorimportance: should knowfreq 40%

answer

  1. same hash fn ⇒ matches share a partition index
  2. 3 IO passes: write partitions, read pair by pair
  3. hybrid keeps partition 0 resident
  4. recursive repartition needs a NEW seed
  5. skew: identical keys can never be split

basics

~20 s

When the build side exceeds memory, the engine hash-partitions both inputs into matching partitions written to temporary storage, then joins one partition pair at a time in memory. That is grace hash join. The hybrid variant keeps the first partition in memory and joins it during partitioning, so part of the work avoids disk entirely.

solid answer

~60 s

**Grace hash join** handles a build input larger than memory by partitioning instead of failing. It hashes the join key of every build row into P partitions and writes them to temporary storage; it then partitions the probe input with **the same** hash function into the same P partitions. Because equal keys always hash to the same partition number, a row in build partition *i* can only match probe rows in probe partition *i*. The engine then processes pairs (build_i, probe_i) independently: load build_i into an in-memory hash table, stream probe_i through it, emit matches. Correctness follows from the partitioning property; cost is one extra write and read of both inputs. **Hybrid hash join** is the practical refinement: while partitioning the build side, keep partition 0 resident in memory rather than writing it out. During the probe partitioning pass, probe rows belonging to partition 0 are joined immediately and never touch disk. With a build side only slightly over budget, most of the join stays in memory. If a partition still exceeds memory — usually key skew — the engine recursively repartitions with a different hash seed, and skew that resists splitting forces a fallback strategy.

code

text · 5 lines
text
->  Hash  (Batches: 32  Original Batches: 8  Memory Usage: 65,536kB  Disk Usage: 2,410MB)
      ->  Seq Scan on line_items  (estimated rows=200,000  actual rows=48,000,000)

Note: Batches grew 8 -> 32 at runtime = the engine repartitioned
      because the estimate was far too low.

go deeper

for a junior

Know that a hash join does not fail when data exceeds memory: it splits both inputs into matching partitions on disk and joins them a pair at a time.

for a middle

Explain the same-hash-function property that guarantees correctness, the extra IO passes involved, and the hybrid optimisation of keeping one partition in memory.

for a senior

Distinguish spills caused by genuine size from spills caused by misestimation, read batch and temp-usage indicators, and handle skew explicitly rather than by adding memory.

for a principal

Weigh memory grants against concurrency and temp-storage bandwidth across the workload, and decide when a spilled hash join is the intended design versus a signal to reshape the query or the model.

## The problem A classic in-memory hash join requires the entire build input to fit in the operator's memory budget. Real systems join tables far larger than memory, and estimates are frequently wrong, so the engine needs a strategy that is correct regardless of size and degrades in cost rather than failing. ## The partitioning insight Everything rests on one property of hashing: **equal keys hash equally**. If you send every build row to partition `h(key) mod P` and every probe row through the *same* function to the same P partitions, then any matching pair must land in the same partition index. That means the single big join decomposes into P independent smaller joins, and no cross-partition comparison is ever needed. Correctness is preserved exactly. ## Grace hash join, step by step 1. **Partition the build input.** Read the build side once, hash each row's key, append it to an output buffer per partition, and flush buffers to temporary storage. Result: P build partition files, each ideally sized to fit in memory. 2. **Partition the probe input.** Read the probe side once with the identical hash function and partition count, producing P probe partition files. 3. **Join pairwise.** For i = 0..P-1: read build partition i into an in-memory hash table, stream probe partition i through it, emit matches, discard the hash table, move on. Cost: each input is read once, written once, and read again — roughly three IO passes instead of one, plus the hashing. That is far cheaper than any quadratic alternative, and it is why data warehouses can join tables vastly larger than RAM. Choosing **P** matters: too few partitions and each is still too big; too many and you pay per-partition buffer memory and produce many small, IO-inefficient files. Engines derive P from the estimated build size and the memory budget, usually targeting partitions comfortably under the budget. ## Hybrid hash join Grace is wasteful when the build side only slightly exceeds memory: it writes *everything* to disk even though most of it could have been held. Hybrid hash join fixes that. During the build partitioning pass, partition 0 (or as many partitions as fit) is kept **resident** in an in-memory hash table instead of being spilled. Then, during the probe partitioning pass, each probe row is inspected: if it belongs to a resident partition it is probed and joined immediately and never written anywhere; otherwise it is written to its partition file. Finally the spilled pairs are processed as in grace. The payoff scales smoothly with the overflow: if 90% of the build fits, roughly 90% of the join is done in memory in a single pass. This graceful degradation is why "hybrid" is what mainstream engines actually implement, even when plans and documentation say "hash join". ## Recursive partitioning and skew After partitioning, a partition may still exceed memory. Two causes: - **Underestimation** — P was chosen too small because the build side was estimated too low. The remedy is **recursive repartitioning**: split the offending partition again using a *different hash seed* (using the same seed would map every row identically and split nothing). Each recursion level adds another IO pass. - **Skew** — one join key value accounts for an enormous share of rows. This is not fixable by repartitioning at all: identical keys hash identically at every level, so they stay together by construction. The engine must fall back to a different strategy for that partition — for example a block nested loop over the oversized group, or in some systems a dedicated skew handling path that treats the hot value separately. Skew handling is a recognised weak point of hash joins, and it is the reason a plan can spend nearly all its time on a single partition. A related refinement some engines use is **role reversal**: if after partitioning it turns out the designated probe partition is much smaller than the build partition, the engine can swap which side builds for that pair. ## Operational signals In a plan or runtime statistics you look for: a batch or partition count greater than one, temporary file volume or bytes spilled, and time concentrated in the hash node. Sharp, non-linear slowdowns after modest data growth are the classic signature of crossing from "fits in memory" to "spills", or of adding one more level of recursion. ## When spilling is fine and when it is a defect It is **fine**, indeed intended, for large batch and analytic joins where the data genuinely exceeds memory. Grace/hybrid partitioning is the designed answer to that, and a stable spilled join with a sensible partition count is a healthy plan. It is a **defect** when unexpected: an OLTP-shaped query spilling, or a report that used to run in memory suddenly spilling, both point at a cardinality misestimate on the build side, missing or stale statistics, or a build side that got wider (selecting unnecessary columns) rather than longer. ## Remedies, in the order worth trying 1. Shrink the build side: push filters down, project only the columns actually needed (narrower rows means a smaller hash table), pre-aggregate. 2. Fix the estimate: refresh statistics, add multicolumn/extended statistics for correlated predicates, avoid predicates the optimizer cannot model. 3. Address skew explicitly: separate the hot key, pre-aggregate it, or split the query for it. 4. Increase the memory budget deliberately, understanding it multiplies across concurrent queries. 5. Consider a different join method if the shape has changed — for example a merge join when both sides are cheaply available in sorted order. ## Summary Grace hash join makes the join size-independent by decomposing it into partition-wise joins that are individually memory-sized; hybrid hash join keeps as much as fits in memory so the extra IO is proportional only to the overflow. Both rely on the fact that equal keys hash to equal partitions — and both are defeated by extreme skew, where equal keys are exactly the problem.

  • Why must the same hash function and partition count be used for both inputs?
    Correctness depends on matching rows landing in the same partition index. If the two sides were partitioned differently, a build row and its matching probe row could end up in partitions that are never compared, silently dropping results. Using one function and one partition count guarantees that equal keys co-locate, so partition-wise joins reproduce the full join exactly.
  • What does the hybrid variant add over plain grace hash join?
    It keeps one or more partitions resident in memory during the partitioning pass instead of writing everything out. Probe rows belonging to a resident partition are joined immediately and never spilled, so the extra IO is proportional only to the portion that did not fit. A build side slightly over budget therefore costs only slightly more than an in-memory join, giving graceful rather than cliff-edge degradation.
  • A single partition is still too large after repartitioning. What is happening and what can be done?
    Almost certainly key skew: one join value holds a huge share of rows, and identical keys hash identically at every recursion level, so no further partitioning can separate them. The engine must fall back to another strategy for that group, such as a block nested loop over it. At the query level the fixes are to handle the hot value separately, pre-aggregate it, or reshape the model so the hot key is not joined at full cardinality.
  • Why must recursive repartitioning use a different hash seed?
    With the same seed, every row in the oversized partition would compute the same partition assignment again and land together, so the split would accomplish nothing. Changing the seed produces an independent distribution over the sub-partitions, which spreads distinct key values that happened to collide at the previous level. It still cannot split a single repeated key value, which is why skew remains unsolved by recursion.

Sorting a warehouse of mail too large for one desk: first drop every letter into pigeonholes by postcode prefix, and drop every delivery slip into the matching pigeonholes. Then work one pigeonhole at a time on the desk — you never have to compare a letter with a slip from another pigeonhole. If one postcode has half the mail, no amount of extra pigeonholes helps.

saying these in an interview costs you the question

  • Believing a hash join fails or errors when the build side exceeds memory
  • Partitioning the two inputs with different hash functions or partition counts
  • Thinking recursive repartitioning solves key skew
  • Reusing the same hash seed when repartitioning
  • Treating every spill as a bug rather than the intended mechanism for large joins

context