skip to content

questions

5

Explain how a database execution engine performs a hash join between two tables, describing what happens in each phase and when the optimizer is likely to choose this join method.

level: middleimportance: must knowfreq 70%

answer

  1. build small side → hash table; probe big side
  2. blocking build, pipelined probe
  3. O(M+N), one pass each
  4. equality only — hashing kills ordering
  5. memory budget ⇒ spill to partitions

basics

~20 s

Phase one (build): read the smaller input fully and hash its join-key values into an in-memory hash table. Phase two (probe): stream the larger input, hash each row's key, look it up, and emit matches. It is chosen for equality joins over large inputs with no useful index.

solid answer

~60 s

A hash join runs in two phases over an **equality** join condition. **Build:** the engine reads one input completely — ideally the smaller one, the *build side* — and inserts each row into an in-memory hash table keyed by the join column(s). This phase is **blocking**: no output can be produced until the build input is fully consumed. **Probe:** the engine streams the other input — the *probe side* — hashing each row's join key and looking it up in the hash table. Every match emits a joined row. This phase is pipelined: rows flow out as they are found. Cost is roughly one scan of each input plus hashing, i.e. O(M + N), versus a nested loop's repeated lookups. That is why optimizers pick it for **large-to-large or large-to-medium equality joins with no selective index**, typical of analytic and reporting queries. It loses to a nested loop when the outer side is tiny and an index on the inner side makes lookups cheap, and to a merge join when both inputs already arrive sorted on the join key.

code

text · 6 lines
text
Hash Join  (cost=... rows=1.2M)
  Hash Cond: (orders.customer_id = customers.id)
  ->  Seq Scan on orders        (rows=12,000,000)      <- probe side
  ->  Hash  (rows=48,000  Batches: 1  Memory: 6,300kB)
        ->  Seq Scan on customers                       <- build side
              Filter: (region = 'EU')

go deeper

for a junior

Describe the two phases in plain terms — build a lookup table from the small side, then stream the big side through it — and note that it only works for equality conditions.

for a middle

Add the cost comparison against nested loop and merge join, the blocking build versus pipelined probe, and that memory bounds the hash table.

for a senior

Discuss build-side selection from estimates, spilling behaviour, skew from duplicate keys, and how to read estimated-versus-actual rows in a plan to diagnose a bad hash join.

for a principal

Frame join-method choice as an optimizer cost decision under estimation uncertainty, and discuss memory budgeting and admission control when many concurrent queries each want hash memory.

## The problem a join algorithm solves Given two row sources and a predicate linking them, produce the matching combinations. The relational model says nothing about *how*; the execution engine has a small family of physical algorithms, and the optimizer picks among them by estimated cost. Hash join is the workhorse for large equality joins. ## Phase 1 — Build The engine chooses one input as the **build side** and reads it to completion. For every row it computes a hash of the join key and stores the row (or the needed columns) in a hash table bucket. Duplicate keys are chained in the same bucket, which is how many-to-many matches are handled. Two important properties: - **The build is blocking.** Not a single output row can be produced until the build input is exhausted, because a later build row might match the first probe row. In a plan tree this shows up as a stall: time-to-first-row includes the whole build. - **The build side should be the smaller input**, because its size determines the hash table's memory footprint. Engines pick this from estimated cardinality, not table size — a heavily filtered big table can legitimately be the build side. ## Phase 2 — Probe The engine streams the other input. For each row it hashes the join key, finds the bucket, and compares keys within it (hashing is not a proof of equality — collisions and different values landing in one bucket both require a real comparison). Each true match emits a joined row immediately, so the probe phase is **pipelined**: downstream operators start receiving rows as soon as probing begins. If the probe side is empty, the work of building was wasted; if the build side is empty, an inner join can short-circuit entirely. ## Why the cost model likes it The classic in-memory hash join reads each input once: cost ≈ scan(build) + scan(probe) + hashing + lookups, i.e. **O(M + N)**. Compare: - **Nested loop** without an index is O(M × N) — fine when the outer side has a handful of rows, catastrophic otherwise. With a selective index on the inner side, it becomes O(M × log N)-ish and is unbeatable for small outer inputs, which is why OLTP point queries use it. - **Merge join** is O(M log M + N log N) if sorting is required, but O(M + N) if both inputs already arrive sorted — for instance from ordered index scans or a previous sort-preserving operator. So the decision rule an interviewer wants to hear: **hash join for large equality joins with no useful index and no free ordering; nested loop for small outer + indexed inner; merge join when the sort order is already there or is needed anyway downstream.** ## Handedness and semantics Because the two sides play different roles, hash join is not symmetric in implementation even though the logical join may be. Outer joins constrain which side can build: to emit unmatched rows from a side, the engine must be able to detect "no match" for that side. A left outer join is naturally implemented by probing with the preserved (left) side and marking matches, or by building on the preserved side and scanning the table afterwards for unmatched entries. Semi-joins (EXISTS-shaped) and anti-joins (NOT EXISTS-shaped) are also implemented as hash joins with early-exit or unmatched-emit behaviour, which is a common follow-up. ## Memory is the defining risk The hash table lives in a bounded memory budget. If the build side is bigger than the budget, the engine must **partition and spill** — the grace/hybrid hash join family — writing partitions of both inputs to temporary storage and joining them pair by pair. This works correctly but costs extra IO, and it depends on the optimizer's row estimate for the build side. An underestimate is one of the most common causes of a plan that looks fine and runs terribly, because the engine sized memory for a small table and met a large one. ## Equality only A hash lookup can only answer "which rows have exactly this key?". It cannot answer "which rows have a key less than this" or "within this range", because hashing destroys ordering. Therefore a hash join can only service **equi-join** predicates. A join whose only condition is an inequality or a range must use a nested loop or a specialized range algorithm; a merge join can handle sorted-order-based band predicates in some engines. In practice, engines apply the equality parts of a predicate in the hash join and evaluate any remaining non-equality condition as a residual filter on matched pairs. ## Reading it in a plan A hash join in a plan sketch typically shows the build subtree and the probe subtree as its two children, with a note about hash table size, batches/partitions, and whether it spilled. Two things to check when performance is bad: (1) was the build side estimated correctly — compare estimated vs actual rows; (2) did it spill — a non-zero batch/partition count or temporary IO indicates the memory budget was exceeded. ## Summary for the interview Build a hash table on the smaller input, probe with the larger, emit matches; one pass over each input; blocking build, pipelined probe; equality predicates only; memory-bound, with partitioning and spilling as the fallback; the right choice for big analytic joins and the wrong choice when a tiny outer input plus an index makes a nested loop nearly free.

  • Which phase is blocking, and why does that matter?
    The build phase blocks: no output row can be emitted until the build input is fully consumed, because a build row read later could match the first probe row. This makes time-to-first-row poor, which matters for queries that only need the first few rows or feed an interactive cursor. The probe phase, by contrast, is pipelined and streams matches to the parent operator as it finds them.
  • When would a nested loop beat a hash join for the same query?
    When the outer input is small and there is a selective index on the inner side's join column. Then each outer row costs an index lookup instead of a full scan, so total work is a few lookups rather than reading both inputs completely. This is the normal OLTP case — fetch one order, look up its customer — where building a hash table over an entire table would be pure waste.
  • How does a hash join handle duplicate join keys?
    Rows with equal keys land in the same bucket and are chained together, so a probe row matching that key emits one output row per chained build row. This is how many-to-many joins are produced correctly. A large duplicate group is also a skew risk: one bucket becomes very long, and probes hitting it do far more comparison work than average.

Checking wedding guests: first you write every invited name onto index cards sorted into labelled boxes by first letter (build). Then, as people arrive, you go straight to the right box and look (probe). You cannot use the boxes to answer "everyone whose name sorts before Miller" — the boxes only answer exact-name questions.

saying these in an interview costs you the question

  • Saying the hash table is built on the larger input
  • Claiming a hash join can serve range or inequality join predicates
  • Believing the probe side is fully materialized in memory as well
  • Saying it produces sorted output — hash join output has no guaranteed order
  • Describing it as always faster than a nested loop, ignoring the small-outer-plus-index case

context

open as a page

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%

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.

open as a page

A query joins two large tables on the condition that one table's timestamp falls between two columns of the other, and the plan will not use a hash join no matter how the query is written. Explain why a hash join can only implement equality join conditions.

level: middleimportance: should knowfreq 45%

basics

~20 s

A hash function maps a value to a bucket and destroys ordering and locality: values that are close numerically land in unrelated buckets. So a hash table can answer "which rows have exactly this key?" but never "which rows are less than, or within a range of, this key." Range joins need other algorithms.

open as a page

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%

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.

open as a page

For a system that runs both short transactional queries and large reporting queries against the same relational engine, how would you decide how much working memory to allow hash join operators, and what are the consequences of setting that budget too high or too low?

level: principalimportance: nice to knowfreq 28%

basics

~20 s

Budget from total memory minus the buffer cache, divided by realistic peak concurrency and the number of memory-hungry operators per query — not per query in isolation. Too low means constant spilling and IO; too high means a few concurrent reports exhaust memory, evict the cache, or crash the process.

open as a page