A grouped run over a 4 GB table exhausts a 32 GB machine's memory. What do you check about the groups first?
answer
- the maximum, not the mean
- the peak follows the largest group
- a small table is no defence
- distinct keys times per-key state
basics
~20 sThe distribution of group sizes, and specifically its maximum. When the per-key step needs its group present, the peak follows the largest group rather than the table, so one key holding most of the rows can exhaust a machine many times the input's size.
solid answer
~50 sLook at the **size of the largest group**, not the average and not the table. A step that must have its group's values at once pays for the biggest group, so a 4 GB table whose rows are spread over a million keys may run fine while the same table with one key covering most of the rows will not. Two figures answer it: the largest group's row count, and the bytes that group materialises. Then check the second axis, because the largest group is not the only way to run out: the accumulators themselves cost the number of distinct keys multiplied by the per-key state, and a very high-cardinality **grouping key** can exhaust memory with no large group anywhere. Finally ask whether the computation needed to hold its group at all — if it folds, neither figure should be driving the peak and the cause is elsewhere.
go deeper
The figure that matters is the size of the biggest group, not the size of the table; a small input can still produce one very large group.
Explain the arithmetic behind the peak — the largest group's row count multiplied by the bytes materialised per row — and why the average group size hides it entirely.
Work the diagnosis in order: the key's distribution, then which axis is large, then whether the computation needed to hold anything, then the cheap remedies before the expensive ones.
Own the standing question: whether a dominant key is data or a defect, and what a scheduled job is allowed to do when the group-size distribution of tomorrow's input is unknown.
## Why the table's size is the wrong number A **grouped operation** — split, apply, combine: one pass that gives every row a key, runs a computation once per key, and reassembles the answers — is not sized by its input when the per-key computation needs that key's rows present together. It is sized by the largest of those groups. The input can be comfortably small and the run can still fail, because the run never had to hold the table; it had to hold one group, and one group happened to be enormous. This is why "the table is 4 GB and the machine has 32 GB" is not a reassurance. The arithmetic that matters is the biggest group's row count multiplied by the bytes materialised per row of it, plus whatever the execution model was already holding. ## The first two figures to obtain Before touching the code, get the shape of the **grouping key**: 1. **The number of distinct keys.** This sets the accumulator bill and tells you whether the grouping is as coarse or as fine as you thought. 2. **The largest group's row count** — not the average, which a single dominant key barely moves. A distribution where the mean group holds forty rows and one group holds sixty million is entirely ordinary in real data, and the mean tells you nothing about it. The common shapes behind a dominant key are worth recognising on sight: a placeholder value standing in for "unknown", a default account or tenant that everything unattributed was assigned to, a key that is coarser than intended, or a single genuinely huge member of an otherwise ordinary population. ## Two axes, not one The largest group is the first thing to look at, and it is not the only way a grouped run exhausts memory. There are two independent axes: | Axis | What it costs | When it dominates | |---|---|---| | the largest group | its row count x the bytes materialised per row | the per-key step must have its group present | | the number of distinct keys | keys x per-key state | a very fine key, even with every group tiny | A run can fail on either. A key with hundreds of millions of distinct values will exhaust a machine while folding a running total, with no group larger than a handful of rows. So the diagnostic order is: establish which of the two axes is large, and only then ask about the computation. ## Whether the computation had to hold anything If the per-key step folds — a fixed-size state, an update taking the state plus one record, a finish step — then group size should not be driving the peak at all, and a failure points somewhere else: - the input was resident because of **eager evaluation**, where each step runs as written and the input is loaded before the split begins; - a step you did not think of as grouped is receiving the group whole, so the materialisation is happening anyway; - the number of distinct keys, not any group, is the real bill; - the output of the run, rather than the run, is what did not fit. If the step does hold — an exact middle value per group, an exact ordering of a group's values, or a **hand-written per-group body**, a function you supply that the library cannot look inside and that is handed the group as a whole — then the largest group is the bill and it is the thing to attack. ## Remedies, in the order worth trying 1. **Narrow what is materialised.** Drop the columns the computation never reads before the split. This changes no answer and can cut the bill by most of it on a wide table. 2. **Replace the holding computation with a folding one where the question allows it.** This is a decision about what the number means, not a tuning knob: an exact middle value and a running figure are different answers, and the choice belongs to whoever owns the metric. 3. **Handle the dominant key separately.** If one key is pathological — a placeholder, an unattributed bucket — it often should not be in the computation at all, and finding that out is itself the result. 4. **Make the key finer,** where a finer key still answers the question, so that no single group is enormous. 5. **Confirm the key is the one you meant.** A dominant group is quite often a signal that rows were lumped together by a key that lost a distinction you needed. ## What the answer demonstrates The candidate worth hiring reaches for the group-size distribution before reaching for a bigger machine, names the maximum rather than the mean as the figure that matters, and knows that a small table is no defence. The candidate who explains that the table is only 4 GB and therefore the failure must be a bug has the wrong model of what a grouped run holds.
- The same computation ran fine last month on a larger table. What changed?Most likely the distribution of the grouping key rather than the volume. A new placeholder or default value collecting unattributed rows, an upstream change that stopped populating a key, or one member of the population growing sharply will all concentrate rows into one group while leaving the total roughly where it was.
- Can a grouped run exhaust memory with no large group at all?Yes. The accumulators cost the number of distinct keys multiplied by the per-key state, so a key with hundreds of millions of distinct values can exhaust a machine while every group holds a handful of rows. That failure looks identical from the outside and has the opposite remedy — a coarser key, not a finer one.
- Why is the average group size an actively misleading figure here?Because it is the row count divided by the number of keys, and one dominant group barely moves it. A distribution whose mean group holds forty rows can contain a group of sixty million. The maximum is the figure the peak follows, so it is the one to obtain.
saying these in an interview costs you the question
- Argues the run cannot fail because the table fits in memory.
- Looks at the average group size rather than the maximum.
- Treats the largest group as the only way to exhaust memory.
- Asks for a bigger machine before profiling the key's distribution.
- Never asks whether the per-key computation needed its group at all.
- Assumes a dominant key is normal data rather than a possible defect.