Two programs compute the same per-key total, but only one lets the runtime fold before the movement — what differs in how they are written?
answer
- when does the runtime learn the operation
- state the combine, or request the group
- whole group means whole input crosses
- same line count, four orders of magnitude
basics
~20 sOne states the per-key combining operation as part of the grouping step, so partial results can be formed on each producing machine. The other asks for the whole group first and computes afterwards, which obliges the job to deliver every record before anything can be combined.
solid answer
~50 sThe difference is when the combining operation becomes known to the runtime. If the program says "group by this key and combine values this way", the producing side can apply that operation to the records it already holds and send one partial result per key. If the program instead says "give me all the records for this key, and then I will compute over the collection", the computation is only defined once the whole group exists in one place — so every record must cross the redistribution first. Both produce the same number, and on a small input both look identical. At scale the first moves one message per key per producer and the second moves the entire input. This is the single most common reason a grouped aggregation is far slower than it needs to be.
code
python · 10 lines# Shape A: the combining operation is part of the grouping step.
# the producer can fold what it holds before anything is sent
step = group_by(key=lambda r: r.country,
combine=lambda a, b: a + b,
value=lambda r: r.amount)
# Shape B: the whole group is requested, then computed over.
# the runtime is told nothing it could apply on the producing side
step = group_by(key=lambda r: r.country) \
.then(lambda country, records: sum(r.amount for r in records))go deeper
Recall the pairing: a step that says how to combine two values for a key can be folded before the records travel; a step that asks for all the records of a key first cannot. Same result, very different volume.
Explain why the runtime cannot fold the second shape — the per-key computation is defined only over the complete group — and describe what each shape sends across the movement, one partial per key per producer against every record.
Show the diagnosis: compare rows produced against bytes moved, then read the step's source for which shape it is. Also say when the group really is needed, such as an exact median, and plan the movement instead of tuning it.
The leverage is in the surfaces a team writes on. Declarative aggregations obtain the fold by default; a library of group-then-compute helpers spreads the expensive shape across a platform, and that is a review standard rather than a per-job fix.
## Two shapes, one answer, very different wires A grouped aggregation forces a **redistribution** (a movement): every worker sends each record it holds to whichever worker will handle that record's grouping key, so equal keys meet. What an author controls is how much has to cross that movement, and the control is exercised entirely through how the step is written. **Shape A — state the combining operation with the grouping.** The program expresses "for each grouping key, combine the values this way". The runtime now knows, before any record moves, what to do with two values for the same key. So each producing worker can apply that operation to the records it already holds and emit one partial result per key — folding locally before the move. The movement carries partials; the collecting side combines them into the final value. **Shape B — ask for the whole group, then compute.** The program expresses "for each grouping key, hand me every record that carries it", and the per-key computation is written over that collection afterwards. The runtime cannot fold anything, not because it is unwilling but because it has not been told what folding would mean here: the computation is defined only over the complete group. Every record therefore crosses the movement, and the collecting side assembles whole groups before the author's code runs. Both shapes return the same numbers. On a laptop-sized input they are indistinguishable. At a billion records over a few hundred keys they differ by four orders of magnitude of network traffic. ## Why the difference is invisible in review | aspect | fold-friendly shape | group-then-compute shape | |---|---|---| | what the program states | the per-key combining operation | that the whole group is wanted | | what may cross the movement | one partial per key per producer | every record | | where the per-key logic runs | partly on producers, finished on the collector | entirely on the collector | | what the collector must hold | partials for its keys | every record for its keys | | line count in the program | usually the same | usually the same | The last row is why this survives code review: the two programs are the same size, often the same shape, and the expensive one frequently reads more naturally, because "get the records, then do the thing" is how people describe the task out loud. ## A third shape: declarative aggregation Where the step is written as a declarative aggregation over a described dataset rather than as functions over records, a planner can recognise the aggregate and insert the fold itself. That is the surface on which the saving is most often obtained without the author thinking about it at all. It is also the surface where the saving silently disappears when the aggregation is replaced by author-supplied logic over a group, for exactly the reason above. ## What actually varies between engines Be careful asserting one model as the model: - **Whether anything is optimised at all.** A declarative surface may rewrite your step substantially; a program written as per-record functions is executed close to literally, and the oldest lineage in this class rewrites nothing and expects the author to supply the local fold explicitly. - **Where the saving lands.** In a finite job cut at a stage barrier — a line across the job where no downstream worker may compute until every upstream piece of the producing step has finished — the producing side writes one bucket per destination to local disk and the collecting side gathers its bucket from every producer. A fold there shrinks the local write and the later fetch. In a continuous job that hands each record over as it is produced, nothing is written between the halves, so folding means holding records briefly and trading delay for volume; some runtimes offer that, others do not. - **What the collecting side does with a whole group.** Some runtimes stream a group past the author's code, others require the group to be materialised before it can be handed over. The second turns a large group into a memory problem as well as a network one, but you should not assume which behaviour you have. ## Recognising it in a running job The diagnosis is nearly always the same: compare the volume entering the movement with the number of rows the step produces. A step that emits a few hundred rows while moving hundreds of gigabytes has not folded. Then read the step's source and ask a single question: 1. Does the program tell the runtime how to combine two values for one key? 2. Or does it ask for the group and combine afterwards? If the answer is the second and the records-per-key ratio is high, rewriting to the first shape is usually the largest single improvement available to that job — and if the computation genuinely needs the whole group (an exact median, a sort within the group, arbitrary cross-record logic), then the movement is real work and the correct response is to plan for it rather than to tune it away.
- Some per-key computations genuinely need the whole group. How do you tell?Ask whether the result can be built by combining two intermediate results for the same key. A total, a count, a maximum or a running extreme can. An exact median, a sort of the group's records, or logic comparing each record to every other cannot be expressed that way, so the group must be assembled. Then the movement is real work, not a mistake.
- The step is written declaratively, so does the shape still matter?Less, because a planner can recognise a declarative aggregate and insert the fold for you. It starts mattering again the moment author-supplied logic over a whole group replaces the aggregate: the planner then has nothing to insert, and the step quietly reverts to moving every record. The visible symptom is the same either way — traffic proportional to records rather than to keys.
- Why does this mistake survive code review so reliably?Because both programs are about the same length, return the same numbers, and the expensive one reads more naturally — "collect the records for this key, then compute" is how the task is described in words. Nothing in the source says how much will cross the network, so the difference only appears at scale, in a movement whose volume tracks the record count.
saying these in an interview costs you the question
- Believes the runtime always rewrites one shape into the other
- Thinks the two shapes differ in the answer, not the traffic
- Says asking for the whole group is fine because memory is large
- Cannot name a computation that genuinely needs the whole group
- Judges the step by line count rather than by what crosses