skip to content

A job groups a billion records by a near-unique identifier, so why does folding locally before the movement save almost nothing?

level: middleimportance: should knowfreq 52%

answer

  1. the win is a ratio, not a size
  2. records per key, not record count
  3. one partial per key, near-unique keys
  4. identifiers group into groups of one
  5. correct, useless, mildly expensive

basics

~20 s

Folding replaces the records sharing a key with one partial result, so its saving is the records-per-key ratio. When keys are near-unique that ratio is about one: every record becomes its own partial, and the traffic is unchanged while the bookkeeping is added.

solid answer

~50 s

A local fold turns the records a machine holds for one grouping key into a single partial result, so the traffic drops by roughly the records-per-key ratio within each producing piece. With a billion records and, say, nine hundred million distinct identifiers, that ratio is close to one: almost every key appears once per piece, so almost every record crosses anyway, now wrapped as a partial. Worse, you have paid for it — a keyed accumulator table sized by the distinct keys in the piece, plus a lookup for every record. The fold is still *correct*, it just has nothing to fold. The lever is not the record count and not the data size: it is how many records share a key, and the saving fades to nothing as the distinct-key count approaches the record count.

go deeper

for a junior

Remember that folding only helps when several records share a grouping key. Grouping by something unique per record, like an event identifier, gives groups of one, so there is nothing to combine before the records travel.

for a middle

Put numbers on it: records, distinct keys, pieces, and the bound of one partial per key per piece. Then name the two costs left behind when the ratio is near one — a lookup per record and a table sized by the distinct keys.

for a senior

Demonstrate that you check key cardinality before blaming the network, and that you know adding a column to a grouping key is the usual cause of a fold silently ceasing to pay. Say what you would do instead when the ratio is genuinely one.

for a principal

The judgement is where the cardinality came from. Near-unique grouping keys are often a modelling artefact nobody revisited, and the cheapest fix is a coarser grain agreed with the consumer rather than more machines behind an unavoidable movement.

## The ratio that governs everything Folding locally before a movement — combining records that share a grouping key on the machine that produced them, so one partial result per key per worker crosses the network instead of every record — buys exactly one thing: it replaces *many* with *one*. So the size of the win is the size of that "many". Write it as three quantities: - **N** — records entering the step; - **K** — distinct grouping keys; - **P** — producing pieces, where a piece is one contiguous share of the input that a single worker processes on its own. Records crossing the redistribution — the movement in which every worker sends each record to whichever worker handles that record's grouping key — go from `N` to at most `min(N, K x P)`. When `K x P` is far below `N`, the fold is transformative. When `K` approaches `N`, `min(N, K x P)` is just `N`, and nothing happened. | shape of the data | N | K | records per key | outcome of folding | |---|---|---|---|---| | events per country | 1,000,000,000 | 200 | 5,000,000 | traffic collapses by orders of magnitude | | events per city | 1,000,000,000 | 50,000 | 20,000 | still a very large win | | events per customer | 1,000,000,000 | 20,000,000 | 50 | a real but modest win | | events per request identifier | 1,000,000,000 | 900,000,000 | ~1.1 | no win; pure overhead | ## Why the near-unique case actually costs you It is not merely neutral. Three costs are real and one of them can end the job: 1. **A lookup per record.** Every record consults and updates a keyed accumulator table. That is cheap individually and not free across a billion records. 2. **An accumulator table sized by the distinct keys in the piece.** With near-unique keys the table is nearly as large as the piece itself, held in the producing worker's own memory. If that memory is short, the table either has to be flushed early — emitting several partials per key, which is fine for correctness — or it becomes a memory-pressure problem, which is a neighbouring subject and not this one. 3. **Slightly larger records on the wire.** A partial result for a single record carries the accumulator's shape rather than the record's, so in the degenerate case you may move marginally *more* bytes than you would have. So the honest statement is: a fold on near-unique keys is correct, useless, and mildly expensive. ## Recognising the case before you run it The question to ask of a grouping key is not "is it big data" but **how many records share one value**. Useful signals: - **Identifiers are the warning sign.** A per-event identifier, an order line reference, a session token or a synthetic primary key are all near-unique by construction. A dimension — country, device class, day, product category, status — is the opposite. - **Estimate the ratio, not the size.** A billion records over a thousand keys is a million records per key; a billion records over eight hundred million keys is barely more than one. Both are "a billion records". - **Composite keys shred the ratio.** Grouping by country alone might give five million records per key; grouping by country plus minute plus device plus campaign can multiply the key count by six orders of magnitude and quietly destroy the fold you were relying on. Adding a column to a grouping key is the most common way a previously fast job becomes slow. ## What to do instead If the ratio is genuinely near one, the fold is not the lever and you should stop trying to tune it: - **Reduce what each record carries** before the movement, so the same number of records costs fewer bytes. What a movement costs in bytes and processor time is its own subject, but the instinct — do not carry fields the step never reads — applies here. - **Ask whether the grouping key needs to be that fine.** Frequently the near-unique key is an artefact: the consumer only ever looked at a coarser rollup, and grouping one level up restores a large records-per-key ratio. - **Accept the movement.** Sometimes a billion near-unique groups is genuinely the computation, and the right answer is to size the work for a full-volume redistribution rather than to pretend a fold will rescue it. ## What varies between runtimes Whether a fold is even attempted differs. On a declarative surface a planner may insert one and, in some systems, measure its effectiveness on the first records and abandon it if the ratio turns out to be near one; other runtimes apply whatever the program stated, unconditionally; and in a continuous job that hands each record over as it is produced, local folding is an optional behaviour that trades delay for volume and is not offered everywhere. Do not assert that "the engine will notice" — some will, many will not, and on near-unique keys the difference is a wasted table per producer rather than a wrong answer.

  • Adding one column to a grouping key made a previously fast job slow. What happened?
    The extra column multiplied the distinct-key count, so the records-per-key ratio collapsed. A fold that used to replace thousands of records with one now replaces one with one, and the movement carries close to the full record volume. The data volume did not change at all; only the key cardinality did, and that is what the fold's saving is made of.
  • If a fold is useless here, is it harmful enough to remove?
    Usually not worth a code change on its own. The costs are a lookup per record and an accumulator table sized by the distinct keys in the piece. The table is the part that matters: on near-unique keys it approaches the size of the piece, so if producer memory is tight, either bound the table and flush partials early or drop the fold for that step.
  • How would you estimate the ratio before running the job at full scale?
    Count distinct values of the intended grouping key against the record count on a sample, remembering that distinct counts extrapolate badly upward from small samples. What you need is only the order of magnitude: if the sample suggests thousands of keys over billions of records, fold; if it suggests the key is close to an identifier, plan for a full-volume movement instead.

saying these in an interview costs you the question

  • Judges the fold by data volume rather than records per key
  • Assumes any grouped aggregation benefits from folding locally
  • Says folding on unique keys is free because nothing merges
  • Forgets the accumulator table grows with distinct keys per piece
  • Believes every runtime detects a useless fold and drops it