skip to content

Aggregates That Fold and Aggregates That Cannot

A grouped run is sized by its largest group, not by the table: an aggregate that must hold every row of a group at once pays for the biggest one, however few rows there are in total.

on this pageshow

questions

4

Computing the largest order value per customer over a long table, what does the run keep for each customer as it reads?

level: juniorimportance: must knowfreq 58%

answer

  1. what is carried between records
  2. one number per customer, not a list
  3. room tracks distinct keys, not rows
  4. compare, keep the bigger, discard

basics

~20 s

One running value per customer — the largest amount seen so far, replaced when a bigger one arrives. Each record is folded in and then dropped, so the room used tracks the number of customers, not the number of orders.

solid answer

~50 s

The grouped operation here — split, apply, combine: one pass that gives every row a key, runs a computation once per key, and reassembles the answers — carries a **fixed-size accumulator** for each distinct customer. For the largest order value that accumulator is a single number: compare the record in hand against it, keep the bigger, discard the record. No two of a customer's records are ever needed at the same moment, so the rows sharing a key never have to be present together. What does grow is the set of keys: the room needed tracks the number of distinct customers, each carrying its one number. The table's rows may still sit in memory for an unrelated reason — under eager evaluation, where each step runs as written, the input was already loaded before the split began — but this reduction does not require them.

go deeper

for a junior

Be able to say what is carried forward: one running value per customer, replaced whenever a bigger amount arrives. That single number is the whole state of this reduction.

for a middle

Explain the update rule and its consequence — the room needed follows the number of distinct keys, and a customer with a million records costs no more than one with three.

for a senior

Distinguish what must be resident from what happens to be loaded: the accumulators always, the input rows only under eager evaluation or when some step demands a whole group at once.

for a principal

The standing call is which per-key computations a scheduled job may use when nobody controls the shape of tomorrow's input, and what the team pays to find out after the fact.

## The operation, named once The **grouped operation** — split, apply, combine — is one pass that gives every row a key, runs a computation once per key, and reassembles the answers into a result. Here the **grouping key** is the customer identifier and the per-key computation is "the largest order value". The interesting question is not what comes out. It is what has to be held while the run is in flight, because that is what decides whether the run finishes. ## What is carried from one record to the next For a largest-value reduction the carried state is one number per customer, and the update rule is small enough to write out: 1. The first record that introduces a customer sets that customer's stored number to the amount on it. 2. Every later record of that customer is compared against the stored number; if the amount is bigger the stored number is replaced, otherwise nothing changes. 3. When the input is exhausted, each customer's stored number is that customer's answer. Notice what is absent from those three steps: at no point are two of a customer's records needed at the same moment, and at no point is a record consulted again after it has been compared. A reduction shaped like this is a **fold** — a fixed-size state, an update that takes the state plus one more record and yields the new state, and a finish step that turns the state into the answer. ## The room a run of this shape needs | What the run holds while it is in flight | How big it is | What it grows with | |---|---|---| | the running largest value for each customer seen so far | one number each | the number of distinct customers | | the key values themselves, so answers can be attributed | one entry each | the number of distinct customers | | the record currently in hand | one record | nothing | | the earlier records of any one customer | nothing at all | — | The honest statement is therefore: for this reduction the room needed tracks the **number of distinct keys**, not the number of rows and not the size of any one group. A billion orders placed by four hundred customers cost four hundred numbers. A million orders placed by a million customers cost a million numbers — the same reduction, a much larger bill, and the bill has nothing to do with how long any customer's history is. ## What is resident anyway, and why that is a separate question It is tempting to jump from "the reduction does not need the rows" to "a grouped run does not hold the table", and that jump is wrong as often as it is right. Two execution models are in play: - **Eager evaluation**, where each step runs as it is written: the input was read into memory before the split began, so the rows are resident whatever the reduction needs. - **Deferred evaluation**, where the tool records the steps as a plan and only runs them when an answer is demanded: a folding reduction can then be carried out with the accumulators alone, and the rows need never all be present. So the durable statement is about what must be resident, not about what happens to be loaded. The accumulators are always resident. The rows are resident under eager evaluation, or whenever some step demands a whole group at once — and this reduction is not such a step. ## Where this shape stops holding The cheapness above is a property of the computation, not of grouping. It disappears when: - the answer depends on the group's whole distribution, such as the exact middle value of a customer's order amounts — the state then grows with the size of the group rather than staying fixed; - the state grows with the group's variety, such as an exact tally of how many different amounts a customer ever used — one entry per distinct amount; - the per-key work is a **hand-written per-group body**, a function you supply that the library cannot look inside: depending on the surface it may be handed one record at a time, or the group's values in one go, or a table of that group's rows, and only the first of those keeps the fixed-size shape. ## What an interviewer is listening for A candidate who answers "you gather each customer's orders and take the biggest" has described the right output and the wrong machine, and that wrong machine is the one that runs out of memory on a table far smaller than the box. The answer wanted is the state and the update rule: one number, compare, keep the bigger, discard. Being able to say that for one reduction is the entry point to saying which other reductions share the shape and which do not.

  • What changes if you also want each customer's second-largest order value?
    The state per customer grows from one number to two, and the update becomes: compare the incoming amount against the smaller of the two kept, and if it wins, insert it and drop the loser. It is still a size you fixed in advance, so the room per customer stays constant and the customer's other records are still discarded as they pass.
  • Does a customer with a million orders cost more room here than one with three?
    No. Both hold one number. For this reduction the room a run needs follows the number of distinct customers, each carrying its fixed-size state, and is unaffected by how the rows are spread across them. That stops being true as soon as the per-key computation needs its group's values available at once.

saying these in an interview costs you the question

  • Says every customer's orders must be gathered before the largest is known.
  • Thinks the room needed grows with the table's row count.
  • Claims any grouped computation must hold its group's rows.
  • Confuses the size of one group with the number of groups.
  • Cannot name what is carried from one record to the next.
open as a page

Which grouped computations finish with one fixed-size running value per group, and which need the group's values available at once?

level: middleimportance: must knowfreq 55%

basics

~20 s

A computation folds when a fixed-size state plus one more record yields the new state — a group's total, its row count, an extreme. It cannot fold when the answer turns on the group's whole distribution, like an exact middle value.

open as a page

When a grouped step must have its group's values at once, is every column of those rows held or only the ones read?

level: middleimportance: should knowfreq 38%

basics

~20 s

It depends on what the surface hands the step. Some pass only the values of the column being reduced; some materialise a table of the group's rows across every column, including ones the step never reads — and then the group's width is part of the bill.

open as a page

A grouped run over a 4 GB table exhausts a 32 GB machine's memory. What do you check about the groups first?

level: seniorimportance: should knowfreq 45%

basics

~20 s

The 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.

open as a page