skip to content

In an equality join, one key value has two million rows on the left and fifty on the right — what must that key's destination produce?

level: seniorimportance: should knowfreq 58%

answer

  1. count both sides, per key
  2. within a key, every pair is emitted
  3. multiply, do not add
  4. concentration times multiplication, on one worker

basics

~20 s

Two million times fifty is one hundred million output rows, all produced by the single worker that owns that key. A heavy key in a join does not only concentrate input on one destination — it multiplies that input into output there.

solid answer

~50 s

Within a key, a join forms every pair: each of the two million left rows is emitted once against each of the fifty matching right rows. That is 100,000,000 rows out of a destination that received barely two million in. Because a key resolves to one destination, all of that is one worker's work — computing the pairs, holding whatever it needs to form them, and emitting the result. The concentration is a factor of roughly a hundred over an even share; the multiplication makes it a factor of five thousand. Two consequences follow: the piece is oversized both to compute and to write, and any downstream step inherits one enormous piece unless the records are redistributed on a different key. How the pairs are formed — a lookup built over one side while the other streams, or both sides sorted and merged — varies by engine and changes the memory profile, not the row count.

go deeper

for a junior

Recall that a join pairs rows: within one key, each left row is emitted against each matching right row, so repeated keys multiply rather than add.

for a middle

Explain that the pairing for a key happens on the one destination the key resolves to, and compute the product for both sides rather than the sum.

for a senior

Demonstrate that you ask for per-key row counts on both sides before the join, and that you can tell a pure concentration problem from a multiplication one by whether the other side's key is unique.

for a principal

Take the position on where an exploding join gets caught — a contract on key uniqueness upstream, a guard in the job, or an operational alert — and argue which of those the organisation can actually sustain.

## What a join does within one key An equality join is, per key, a cross product. Every record on the left carrying key *k* is paired with every record on the right carrying key *k*. If the left has *L* such rows and the right has *R*, the key emits *L × R* rows. For the overwhelming majority of keys in a healthy dataset that product is tiny — 20 left rows against 3 right rows is 60 — and nobody thinks about it. The join is also a **redistributing step**: a step that cannot be computed from what one worker already holds, so every record is first sent to the worker that owns its key. Which worker that is comes from the **key-to-destination rule**, the function from a key value to the number of the destination that receives it, with the property that the same value always yields the same destination. So both sides' records for key *k* arrive at one destination, and the product for *k* is formed there, by one worker. ## The arithmetic | Key, left × right | Records arriving at the destination | Records produced there | | --- | --- | --- | | a typical key, 20 × 3 | 23 | 60 | | the heavy key, 2,000,000 × 50 | 2,000,050 | 100,000,000 | | the heavy key, 2,000,000 × 1 | 2,000,001 | 2,000,000 | | heavy on both sides, 2,000,000 × 2,000,000 | 4,000,000 | 4,000,000,000,000 | Read the first two rows together. Against a typical key the heavy key brings a hundred thousand times as much input, but produces 1.6 million times as much output. The multiplier on the right-hand side is what turns an awkward piece into one that does not finish. Row three is worth stating explicitly because it is the case people assume: when the other side's key is unique, the product collapses to *L* and the heavy key is a pure concentration problem. Row four is the one that produces a job nobody can wait for. Four trillion rows out of one worker is not slow; it is never. ## Concentration and multiplication are different problems - **Concentration** is about the input: one destination receives a large share of the records. It is set by the key's frequency in the data. - **Multiplication** is about the output: that destination emits the product of the two sides' counts. It is set by the frequencies on **both** sides together. A candidate who only counts input rows will predict a piece a hundred times too large and be confused when the job behaves far worse than that. The counts to look at before a join are the row counts per key on **each** side, not the table sizes. ## What varies between engines - How the pairs are formed at the destination differs. One common approach builds a lookup over one side's rows for the key and streams the other side against it, which means that side's rows for the heavy key must be held somewhere while the pairing runs. Another sorts both sides by key and merges them, which needs no per-key lookup at all and can stream both sides. The memory profile of the heavy key differs sharply between the two; the number of output rows does not change at all. - Whether the runtime intervenes differs too. Where an engine re-decides the part of a plan it has not run yet from the sizes a finished step actually produced, it may treat an unusually large destination specially for a join — for some join shapes. Where the input never ends, no step finishes and there is nothing measured to re-plan from. - What the destination does when it cannot hold what it needs varies as well: writing to local disk and continuing, or failing the piece. Either way, the heavy piece is the one that announces the imbalance first. ## The output does not stay put The destination that produced a hundred million rows is also the one that emits them. If the job writes at that point, it writes one enormous unit of output while its peers write small ones. If another step follows, that step inherits pieces of wildly different size — the imbalance propagates downstream unless the records are redistributed again on a key with a different frequency shape. This is why a heavy key discovered at a join is often blamed on a step several stages later: the later step is where the oversized piece finally fails. ## What an interviewer is listening for That you multiply rather than add. The strong answer states the product immediately, names the destination that will form it, and then asks for the per-key row counts on the other side — because the difference between 2,000,000 × 1 and 2,000,000 × 50 is the difference between an awkward job and one that will not complete. A candidate who says "the join output is about the size of the bigger input" has assumed uniqueness on the other side without checking it.

  • What changes if the heavy key is heavy on both sides?
    The product becomes quadratic. Two million against two million is four trillion pairs from one key, formed and emitted by one worker. Input grew by a factor of two; output grew by a factor of two million. In practice the piece does not complete, and no amount of capacity changes that, because a key resolves to one destination.
  • Does making it an outer join change the arithmetic for the heavy key?
    Not materially. An outer join adds rows only for keys that found no partner, one per unmatched row. A heavy key present on both sides still contributes its full product, and that product dominates everything the preservation rule adds.
  • If the heavy key has no match on the other side at all, is the join cheap for it?
    The output is empty for that key, but the work is not free: its two million rows were still sent to a destination and grouped there before the absence of a partner was discovered. Some engines prune keys that cannot match before they are moved; where that does not happen, you pay the move and get nothing.

saying these in an interview costs you the question

  • The join output is about the size of the larger input, so it is cheap.
  • Only the side being scanned needs its per-key row counts checked.
  • A heavy key in a join is a memory problem, not an output-size problem.
  • Pairing for one key is shared across destinations, so the work is spread.
  • If the heavy key has no match on the other side, none of its rows were moved.