skip to content

Why does an execution engine build the hash table on the smaller of the two join inputs, and what goes wrong at runtime when the optimizer's row estimate for that input turns out to be far too low?

level: seniorimportance: must knowfreq 55%

answer

  1. hash table sized by build side ⇒ pick smallest estimate
  2. smallest after filters, not smallest table
  3. estimate too low ⇒ spill, or sides inverted
  4. est vs actual rows on the build subtree
  5. skew defeats repartitioning

basics

~30 s

The build side's size determines the hash table's memory footprint, so the smaller input keeps it in memory. Choice is made from estimated rows, not actual. A large underestimate produces a hash table far bigger than its memory grant, forcing spilling to disk, or leaves the engine probing a huge table against a tiny one — either way, a plan that looked cheap runs for a very long time.

solid answer

~1 min

The hash table must hold the entire build input, so its size *is* the memory cost of the operator. Choosing the smaller input keeps the table in the memory budget, minimizes hashing work, and gives the best chance of a single-pass join. Note that "smaller" means smaller **estimated result after filters and projections**, not the smaller table — a heavily filtered fact table can legitimately be the build side. The choice is made at plan time from cardinality estimates. When the estimate is badly low — typically from correlated predicates, skewed values, stale statistics, an expression the optimizer cannot model, or an estimate propagated through several joins — two things go wrong: - The hash table exceeds the memory grant, so the engine spills partitions of **both** inputs to temporary storage and rejoins them pass by pass, adding large amounts of IO. - The sides may be inverted: the engine builds over what is actually the huge input and probes with the small one, doing far more work than the mirrored plan would. Diagnose by comparing estimated versus actual rows on the build subtree and checking for spill/batch counts. Fix the estimate (statistics, extended/multicolumn statistics, rewritten predicate) rather than the symptom.

code

text · 7 lines
text
Hash Join  (actual time=... rows=9,400,000)
  Hash Cond: (f.dim_id = d.id)
  ->  Seq Scan on facts f      (rows=9,400,000)
  ->  Hash  (Batches: 64  Disk Usage: 512MB)
        ->  Seq Scan on dims d
              Filter: (city = 'Springfield' AND state = 'IL')
              (estimated rows=500  actual rows=4,800,000)

go deeper

for a junior

Know that the hash table holds the build input, so the smaller input is chosen to keep it in memory, and that a wrong size guess makes the join slow.

for a middle

Add that the choice comes from plan-time estimates after filters, and that exceeding memory causes spilling to temporary storage.

for a senior

Diagnose from estimated-versus-actual rows and spill indicators, name the specific estimation failures (correlation, skew, expressions, stale stats), and fix the estimate rather than the grant.

for a principal

Discuss estimation error as a systemic risk under concurrency — many queries misestimating at once — and design choices that bound it, such as materializing intermediates with known cardinality and provisioning for spilled joins deliberately.

## Why build-side size is the whole cost story A hash join materializes one input entirely — the build side — into a hash table. That table's footprint is roughly (rows × width of the needed columns) plus per-entry overhead for pointers, hash values, and bucket structure; the overhead is not small, so a hash table is typically noticeably larger than the raw data it holds. The probe side, by contrast, is streamed: it costs a scan but essentially no memory. Therefore: - **Memory** — the build side sets the operator's memory demand, and whether the join stays in one pass. - **CPU** — inserts are more expensive than probes, so hashing the smaller input is cheaper in absolute work. - **Robustness** — a small build side gives headroom for estimate error; a build side already near the budget will spill on any underestimate. Hence the rule: build on the input with the smaller **estimated output cardinality after filters**, not the physically smaller table. A 10-billion-row fact table filtered to 3,000 rows is the right build side against an unfiltered 5-million-row dimension. ## Where the estimate comes from and why it breaks The optimizer estimates rows using statistics: per-column histograms, distinct-value counts, most-common-value lists, null fractions. It then combines predicate selectivities, and here is where it goes wrong: - **Correlation.** With `city = 'Springfield' AND state = 'IL'`, independent multiplication of selectivities produces a wildly low estimate, because those columns are not independent. Multicolumn/extended statistics exist precisely for this. - **Skew.** A default histogram may miss that one key value accounts for 40% of the rows, so an estimate based on average frequency is far off for that value. - **Opaque expressions.** Predicates on function results, complex expressions, or values not known at plan time (parameters, subquery results) fall back to fixed default selectivity guesses. - **Stale statistics.** Bulk loads and steady growth leave the optimizer reasoning about a table that no longer exists in that shape. - **Error propagation.** In a multi-join query, an error at the bottom compounds upward; by the third join the estimate can be off by orders of magnitude in either direction. Estimates rarely improve as they propagate. ## What the runtime does when the build side is far bigger than expected **1. Spilling.** Once the hash table exceeds the grant, the engine falls back to partitioning: it splits both inputs into partitions by a second hash of the key, writes them to temporary storage, and then joins one partition pair at a time in memory. Correct, but it converts an in-memory operation into a multi-pass IO-heavy one. Runtime does not degrade gracefully — the shift from "fits" to "spills" is a step change, and a plan that ran in seconds yesterday can take many minutes today. **2. Repeated repartitioning.** If a single partition still does not fit — usually because of key skew concentrating rows into one partition — the engine may repartition recursively. Extreme skew (one key with a huge fraction of rows) defeats hash partitioning entirely, since all those rows hash to the same place and cannot be split further by the same function. **3. Inverted sides.** A worse outcome than spilling: the estimate flips which input looks smaller, so the engine builds on the genuinely huge input and probes with the small one. Now the memory demand is maximal and the cheap streaming side is the tiny one, wasting the algorithm's entire advantage. **4. Concurrency pressure.** Memory grants are per-operator and per-query. Several concurrent queries all under-estimating and all over-allocating (or all spilling) will exhaust the machine's memory or its temporary storage bandwidth, so a single bad estimate can degrade unrelated queries. ## Diagnosing it The single most valuable signal is **estimated versus actual rows** on the build subtree. Any plan display that reports both makes the diagnosis nearly immediate: an estimate of 500 with an actual of 5,000,000 explains everything downstream. Also look for reported spill indicators — batch/partition counts above one, temporary file or IO volume, or a memory-spill wait — and for actual time concentrated in the hash node. ## Fixing it properly 1. **Repair statistics** — refresh them, raise histogram resolution on skewed columns, and create multicolumn/extended statistics for correlated predicates. This addresses the cause. 2. **Rewrite the predicate** so the optimizer can estimate it — avoid wrapping columns in functions, avoid type mismatches that force conversions, and expose literal-comparable columns. 3. **Reduce the build side** — push filters and projections down so fewer, narrower rows are hashed. Selecting only the needed columns genuinely shrinks the hash table. 4. **Reshape the query** — pre-aggregate or materialize an intermediate result with known cardinality so the estimate cannot drift through several joins. 5. **Grant more memory** only as a deliberate, measured decision — it treats the symptom and competes with concurrency. 6. **Accept spilling by design** for genuinely huge joins: a correctly partitioned, spilled hash join is the normal way to join data larger than memory, and it is a healthy plan when the sizes are known and stable. ## The judgment an interviewer is listening for That you know the build side is chosen by *estimate*, that estimation error is the dominant failure mode of hash joins in production, that the symptom is a step change from fast to very slow rather than gradual decay, and that the durable fix is to make the estimate right rather than to enlarge the memory grant.

  • Name the estimation errors that most often cause this.
    Correlated predicates whose selectivities are multiplied as if independent, value skew that a coarse histogram does not capture, predicates over expressions or function results that fall back to default selectivity, stale statistics after bulk loads, and error propagation through multi-join plans where a small mistake at the bottom compounds upward. Multicolumn or extended statistics address the correlation case directly.
  • Why is the performance cliff so sharp when a hash join spills?
    In-memory probing is a pointer chase; spilling converts the operator into partitioning both inputs to temporary storage and rejoining pass by pass, so the cost model changes category from memory access to IO. There is no gradual middle ground: either the hash table fits the grant or the whole operator switches strategy. That is why the same query can go from seconds to many minutes after a modest data-volume change.
  • How does key skew interact with spilling?
    Hash partitioning splits data by another hash of the key, so all rows sharing one very frequent key land in the same partition no matter how many partitions are created. If that single partition still exceeds memory, recursive repartitioning cannot help, and the engine must fall back to a slower strategy for that partition. Skew is therefore worse than mere size, and it is usually addressed by separating or pre-aggregating the hot key.
  • Is spilling always a defect?
    No. Joining inputs larger than available memory is exactly what the grace/hybrid partitioning design exists for, and a stable, correctly partitioned spilled join is a healthy plan for large batch work. It is a defect when it is unexpected — that is, when the build side spilled because the estimate was wrong and the query was provisioned as an in-memory join.

saying these in an interview costs you the question

  • Saying the engine picks the smaller table rather than the smaller estimated result after filters
  • Assuming the build side is chosen at runtime from actual sizes
  • Treating a spilled hash join as always broken rather than sometimes intended
  • Fixing repeated spills by permanently raising the memory grant without addressing the estimate
  • Ignoring estimated-versus-actual rows when diagnosing a slow hash join

context