A grouped total over 80 million rows with 50 million distinct keys is slow — which phase is costing you?
answer
- count the distinct keys first
- trivial reduction, expensive keying
- the middle phase grows when you write it
- ask what the supplied body receives
basics
~20 sThe split. With a near-unique key, organising tens of millions of key values dominates: each apply touches only one or two rows, and the combine emits almost as many result rows as there were inputs.
solid answer
~60 sTwo variables set where the cost lands: how many distinct keys there are, and what the apply does. With a **near-unique key** the split dominates — one key evaluation per row, then sorting or hashing fifty million values, then a combine emitting nearly one result row per input row; the arithmetic itself is trivial. With **few keys and a reduction the library implements**, the split is a cheap bucketing and the apply is one linear fold, so the whole operation is close to a single sweep over the column. A **hand-written per-group body** moves cost into the middle, and how much depends on what the body receives: a body handed one record at a time pays a boundary crossing out of the library's own execution for every record, while a body handed the whole group's column at once runs over that column in one go and pays nothing per record. So the first question of a slow grouped job is: how many distinct keys, and is the apply something the library implements itself?
go deeper
Know that a grouped operation has a cost before any arithmetic happens: every row's key has to be worked out and the distinct values organised.
Explain how the number of distinct keys changes the shape of the cost, and why a result nearly as large as the input means the operation condensed nothing.
Diagnose in order: count distinct keys, check whether the apply is a library reduction, and find out what a supplied body is actually handed before changing anything.
Challenge the design when a key is near-unique: a computation per identifier is a row-level computation, and paying a grouped operation's keying cost for it is a standing waste.
## Two variables, not one People reach for "the aggregation is slow" as though the arithmetic were the cost. It rarely is. Where a grouped operation spends its time is set by two things: 1. **The number of distinct key values** — this sizes the split's bookkeeping and the combine's output. 2. **What the apply is** — a reduction the library implements, or a body it cannot look inside. Everything below follows from those two. | distinct keys | what the apply is | dominant phase | why | | --- | --- | --- | --- | | few | a library reduction | apply, and it is cheap | one linear fold over the column; near one sweep total | | few | a supplied body | apply | one invocation per key, but each sees many rows | | near-unique | a library reduction | split, then combine | organising tens of millions of values; output nearly as large as input | | near-unique | a supplied body | split and apply together | the worst case: expensive keying plus millions of invocations | ## Near-unique keys: the split dominates At fifty million distinct keys over eighty million rows the arithmetic is nothing — most groups contain one or two rows. What costs: - **One key evaluation per row.** Eighty million evaluations, before any grouping exists. - **Organising fifty million distinct values**, either by ordering the key column or by hashing into buckets. Either way this is the real work, and the structure holding it is comparable in size to the column you grouped on. - **A combine emitting nearly one output row per input row.** The result is barely smaller than the table, so nothing has been saved by aggregating. The diagnostic conclusion is often that the grouping is the wrong operation: a key that is nearly unique is an identifier, and a computation per identifier is a row-level computation wearing a grouped operation's clothes. ## Few keys and a library reduction: nearly free The opposite end is the case people wrongly imagine is always slow. With a dozen groups: - The split's bookkeeping is a dozen entries plus one bucket assignment per row. - The apply is a single fold per group over values read from the original storage. - The combine emits a dozen rows. This is close to one pass over the column, and adding more rows costs linearly with no change in shape. ## A supplied body moves the cost into the middle — but not always by the same amount This is the claim most often stated too strongly. "A function you hand in is called once per record, so it is always slow" is true of one calling convention and false of another, and the two look nearly identical from outside: - Where the surface hands your body **one record at a time**, you pay a crossing out of the library's own execution for every record. At eighty million rows that dominates everything else in the operation. - Where the surface hands your body **the whole group's column at once**, the body runs over that column in one go — expressed once for the whole column rather than once per row — and the per-record penalty never appears. So the question to ask about a supplied body is never "is it hand-written?" but **"what is it handed?"** A secondary cost applies in the whole-table case: gathering every column of the group side by side is a materialisation, which is memory traffic the reduction path never pays. ## Reading your own job A workable order of investigation: 1. **Count the distinct key values.** If that number is close to the row count, stop — the split is your cost, and the design is probably wrong. 2. **Check whether the apply is a reduction the library implements.** If it is, and the key cardinality is low, the job should already be near one sweep; something else is wrong. 3. **If a body is supplied, find out what it receives.** One record at a time and a whole column are different cost classes under a similar-looking surface. 4. **Look at the combine's output size.** A result nearly as large as the input means the operation is not condensing anything. ## The trap in "the split does no arithmetic" Because the split computes no answer, it reads as preparation rather than as work, and estimates routinely leave it out. It is the phase whose cost grows with something people do not think to measure — the number of distinct keys — and it is the most common answer to "why is this grouped job slow".
- Two grouped jobs read the same table; one produces 12 groups and one produces 50 million. What differs?Where the time goes. With twelve groups the split is a cheap bucketing and the apply is a linear fold, so the job is near one sweep. With fifty million, organising the key values dominates and the combine emits nearly one output row per input row, so almost nothing is condensed.
- Does using fewer grouping key columns always make the operation cheaper?Dropping a key can only merge groups, so the group count falls or stays the same, which usually cuts the split's and the combine's cost. It also means each apply sees more rows, so the effect on the middle phase depends on what that apply does with them.
- Why is 'the aggregation function is slow' usually the wrong first hypothesis?Because a library-implemented reduction is a linear fold over values already in storage, which is the cheapest part of the operation. The expensive parts are evaluating and organising the keys, and any per-record crossing a supplied body forces.
saying these in an interview costs you the question
- Assumes the reduction is always where a slow grouped job spends its time.
- Says a supplied per-group body is always slow, whatever it is handed.
- Ignores the number of distinct keys when reasoning about grouped cost.
- Treats the split as free because no arithmetic has happened there.
- Reaches for a bigger machine before counting the distinct key values.