You spread a heavy grouping key over 200 artificial suffix values and attach the suffix to every key - what does the width buy, and what does attaching it everywhere cost?
answer
- the width is measured, not guessed
- heavy key records over a typical destination
- parity is the target, not zero
- suffixing every key multiplies the middle
- spread only the measured heavy values
basics
~20 sWidth divides the heavy key's records by that factor and buys nothing once its share drops below an ordinary piece's. Attaching the suffix to every key multiplies the intermediate result by 200, and the second pass must move and read all of it.
solid answer
~50 sThe width is a measurement, not a preference. Pick it from the ratio you already have to measure to diagnose the problem: the heavy value's record count divided by a typical destination's. Once the spread share is comparable to an ordinary one, that destination is no longer the slowest and widening further changes nothing about the run. The cost sits in the middle of the job. The first pass emits one partial per `(key, suffix)`, so a suffix on every key makes the intermediate result `distinct keys x 200` rows for the second redistributing step to move and read. With tens of millions of distinct keys that is billions of partials, comfortably worse than the imbalance. Attach the suffix only to values measured heavy, give everything else a constant suffix, and the middle grows by `heavy keys x 200` instead.
go deeper
Know that the suffix has a width, that it decides how many destinations one key's records reach, and that a wider spread is not automatically a better one.
Explain the arithmetic both ways: the busiest share falls by the width, and the number of rows between the two passes rises by it, multiplied by however many keys carry a suffix.
Derive the width from measured record counts, restrict the suffix to a measured list of heavy values, and say what happens to that list and that width as the distribution moves.
Decide who owns the measurement over time: a rewrite carrying hand-tuned constants is a standing maintenance cost, and it is worth stating which jobs justify it and which should absorb the slow run.
## What the width is for Spreading a heavy key means appending an integer `s` in `0..W-1` to a grouping key whose one value covers a large share of the records, so those records reach several destinations instead of the single one the key-to-destination rule would send them to, and then folding the partial results under the original key in a second pass. `W` is the width, and it is the only free parameter in the rewrite. Its effect is arithmetic. The heavy value's records are divided over up to `W` destinations, so the busiest destination's load falls to roughly `records under the heavy key / W`, plus whatever ordinary keys happen to land there. The target is parity, not zero: > `W` is about `records under the heavy key / records in a typical destination`, rounded up. If a value holds 400 million records and a typical destination holds 4 million, a width near 100 brings it into line. Past that point the heavy value is no longer the slowest thing in the step, the run cost is set by something else, and further widening is pure overhead. This is why the width follows the per-destination record-count measurement rather than preceding it. ## Where the width is paid for The first pass emits one partial per `(key, suffix)` pair that actually occurs. Three costs follow: - **Groups held during the first pass.** Each destination tracks more distinct groups, which is memory on the gathering side; a heavy piece is the one that spills to local disk or dies first, and that is how the pressure announces itself. - **Rows the second redistribution moves.** This is the dominant one. Every partial is a record that has to be sent and read again. - **Combining work.** Trivial per key, but multiplied by the number of partials. Now apply the suffix to every key. A job with 40 million distinct keys and `W = 200` produces up to 8 billion partial rows where the plain grouping produced 40 million results. The second pass is no longer a cheap fold over summaries; it is a bigger redistribution than the one you were trying to balance. A width that is right for the heavy value is catastrophic for the other 39,999,999 keys, and that asymmetry is the whole reason known-heavy values are handled apart. ## Handling the heavy values apart 1. **Measure.** Count records per key - itself a grouping, but a cheap and decomposable one - or count over a sample large enough that a value covering a large share cannot hide in it. 2. **Build a small list.** Keep the values above a chosen threshold. There are normally few of them, small enough for every worker to hold a copy. 3. **Suffix conditionally.** In the first pass, attach `s` in `0..W-1` only when the key is on the list, and a constant otherwise. Ordinary keys then emit exactly one partial each. 4. **Fold uniformly.** The second pass drops the suffix and groups by the original key. It needs no branch: a key with one partial folds to itself. The middle of the job is then `heavy keys x W + other keys` rows, which is the original size plus a rounding error. A refinement worth knowing is per-value widths - each heavy value spread in proportion to its own count - so that one value covering half the input and one covering a twentieth do not get the same treatment. ## Keeping it honest - **The list goes stale.** Key frequencies move, and a width tuned to last quarter's distribution can be both too small for the new heavy value and dead weight for the old one. Either re-derive the list and the width from a counting pass inside the job, or schedule a review of the measurement rather than of the code. - **The suffix should survive recomputation.** Where a runtime recovers a lost piece by recomputing it from its inputs, a suffix drawn at random per record makes the replacement piece place records differently from the one it replaces. A suffix derived from a varied secondary field is stable; round-robin within a worker is the usual compromise. - **Some of this may already be done for you, and some of it never will be.** Some runtimes re-decide the part of a plan they have not run yet using the sizes a finished step actually produced, and will divide an oversized share on their own - for certain operators only. Others do nothing of the kind. A job over an endless input gets very little of that help, because no step ever finishes to be measured, and there the extra groups are retained for the life of the job rather than for the length of one run, which turns the width from a one-run cost into a standing one. Establish which of these you are on before hand-writing the rewrite.
- How do you find the values worth spreading without paying for a full count?Count records per key over a sample. A value covering a large share of the input cannot hide in a sample of any reasonable size, and the counting itself decomposes into partial counts, so it does not inherit the imbalance you are diagnosing. Exact counts matter only for setting the width, and the width tolerates being approximate.
- Should every heavy value get the same width?Not necessarily. Width in proportion to each value's record count keeps every spread share near a typical destination's load, so a value covering half the input gets a wide spread and a merely large one gets a narrow one. It costs a lookup of the per-value width in the first pass and nothing in the second.
- What tells you the width was too small?The same signature as before the rewrite, scaled down: one destination in the first pass still reads many times the input of its peers and still sets the step's duration. Compare the busiest share against the median share after the change - if the ratio is unchanged, the spread did not reach the value that matters.
saying these in an interview costs you the question
- Use the largest width you can; more spreading is always safer
- The suffix costs nothing for keys that are not heavy
- Once tuned, the width never needs revisiting
- Measuring records per key is too expensive to bother with
- A random suffix can be redrawn on recomputation without consequence
- If the job got faster at all, the width was right