In a step that gathers every record with one key onto a single worker, why does doubling the cluster not shorten the commonest key's work?
answer
- the run costs its slowest piece
- who decides where a key goes
- same key always, same destination
- one value cannot span two destinations
basics
~20 sA key-to-destination rule sends every record carrying the same key to one destination, and that mapping is fixed. Extra machines add destinations, not a way to split one key, so the commonest key still runs on a single worker.
solid answer
~50 sA grouped aggregate or an equality join cannot be computed from what one worker already holds, so the engine inserts a redistributing step: every record is first sent to the worker that owns its key. Ownership comes from a key-to-destination rule, a function from key value to destination number, and its defining property is that the same key value always yields the same destination. That determinism is what makes the records for a key meet, and it is also the limit: one key is one destination's work. If a single value covers 40% of the records, doubling the cluster buys idle capacity, not a smaller heaviest piece — a run costs what its slowest piece costs. Only two things move that floor: changing the key values, or changing the step so co-location by that key is no longer required.
go deeper
Recall the shape: records sharing a key must meet on one worker, so the slowest piece is the one holding the commonest key. More machines do not divide that piece.
Explain the mechanics: the key-to-destination rule is deterministic, so one key value maps to exactly one destination, and a larger cluster only changes how many destinations can run at once.
Show that you state the limit before proposing anything. With the key and the step fixed, the heaviest key is the floor under the run, and capacity spend buys nothing against it.
Frame it as where the fix belongs — in the upstream data that produced the key, in the job's logic, or in accepting the runtime — and price each against engineer time rather than machine time.
## The step that forces records together Most work in a distributed job divides trivially. A filter, a projection, a derivation from the fields of the same record — each can be computed from what one worker already holds. The input is cut into **pieces** (one share of a step's input that one worker processes independently of the others; how many pieces there are sets how much of the cluster works at once), and every piece runs without talking to any other. A few operations cannot work that way. A grouped count, a per-key sum, an equality join: the answer for one key depends on records that may sit on any worker in the cluster. The engine handles this with 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. The industry's word for this family of steps is a *shuffle*. The workers that emit records into it are the **producing side**; the workers that receive them, already grouped by key, are the **gathering side**. ## The rule that decides who owns a key Ownership is decided by the **key-to-destination rule**: the function that turns a record's key into the number of the destination that receives it. Its essential property is determinism — the same key value yields the same destination number, on every producing worker, for every record. That is not an incidental detail. Determinism is the whole reason the records for a key meet at all: two producing workers that never communicate both compute the same destination for the same key, and the records converge there. And the same property is the limit. If one key value always yields one destination number, the records for that key are by construction one destination's work. There is no way to place half of them elsewhere and still satisfy what the step promised. ## Why capacity does not help Now add a **heavy key**: a single key value that covers a large share of all the records — a default account, a marketplace's busiest seller, a placeholder meaning "unknown". Take a billion records, four hundred pieces on the gathering side, and one key holding four hundred million of them. Divide the remaining six hundred million evenly and each of the other pieces handles about 1.5 million records. One piece handles four hundred million. That is a ratio of roughly 265 to 1, and the job is not finished until the last piece is. - Adding worker machines raises how many destinations can run **at the same time**. It does not change which destination a key resolves to. - Raising the number of pieces divides the *other* keys more finely: each destination receives fewer distinct key values. The heavy key is a single value and still arrives whole. - Giving that worker a larger machine can let an oversized piece survive instead of failing or writing to local disk, but it does not remove the imbalance — the run is still gated by that piece. - Re-running changes nothing, because the rule is deterministic: the same key lands on the same destination number every time. | What you add | What it changes | Effect on the heavy key | | --- | --- | --- | | more worker machines | more destinations can run at once | none; the key still resolves to one | | more pieces on the gathering side | the other keys spread more finely | none; a single value cannot be divided | | a larger machine for that worker | headroom for one oversized piece | it may survive rather than fail, but still gates the run | | a re-run of the same job | nothing | the mapping is deterministic | ## The two axes that are left Hold the step's requirement fixed and only two things can actually move: 1. **The key values** — so the records no longer share a single value and therefore no longer share a single destination. 2. **The step** — so co-location by that key is no longer required at all, or so that far fewer records reach the gathering side. Every real fix for a heavy key acts on one of those two axes. Anything that only adds capacity acts on neither, which is why "give it more machines" is the answer an interviewer is waiting to hear rejected. ## What varies between engines - Whether the runtime can intervene by itself varies. Where an engine re-decides the part of a plan it has not yet run, using the sizes a finished step actually produced rather than the estimate it started with, it may treat a very large destination specially — for some operators only. Where the input never ends, no step ever finishes, so there is little measured evidence and usually no such help. - Whether redistributed records are written down and fetched afterwards, or pushed across the network as they are produced, differs by engine lineage. It changes *when* the imbalance becomes visible and what it costs in storage; it does not change that one key is one destination's work. ## What an interviewer is listening for That you name the constraint before proposing anything: all records for a key must meet, one key resolves to one destination, so the heaviest key sets a floor under the run that no amount of hardware touches. A candidate who reaches for cluster size first has not understood which quantity the runtime is actually bounded by.
- Does the same limit apply to a step that only filters or reshapes each record on its own?No. A step computable from what a worker already holds — a filter, a per-record derivation — needs no co-location, so its work divides with the number of pieces whatever the key frequencies look like. The limit appears only where a step's contract is that all records for a key be present in one place.
- Can the runtime notice the heavy piece and divide it for you?Some engines can, for some operators. Where a runtime re-decides the part of a plan it has not run yet from the sizes a finished step actually produced, it may handle an unusually large destination specially. Whether that happens depends on the engine and the operation, and in a job over an endless input no step ever finishes, so there is usually nothing to measure and no such help.
- If one key covers 40% of the records, what is the best speedup any number of machines can give?The run cannot finish faster than that one piece. Even with the remaining 60% spread perfectly and instantly, the floor is the time one worker needs for 40% of the data — so the achievable speedup over a single machine is at most about 2.5x, no matter how large the cluster is.
A mail room sorts letters into pigeonholes by street name, because every letter for a street has to be bundled together before it goes out. One street carries 40% of the town's post. Building more pigeonholes sorts the other streets more finely and leaves the clerk on that street with exactly the same mountain. The only ways out are to write the addresses differently or to stop bundling by street.
saying these in an interview costs you the question
- Adding worker machines will make the heavy key finish sooner.
- Raising the number of pieces divides the records of one key.
- The runtime always detects a heavy piece and splits it automatically.
- A bigger machine for that worker removes the imbalance.
- Every step in a job forces records with the same key together.