skip to content

A per-key total runs through a step your runtime will not fold locally, so every record crosses the network — how do you fold it by hand, and what does that cost?

level: seniorimportance: should knowfreq 44%

answer

  1. two steps where the intent had one
  2. accumulate inside the piece, then exchange
  3. bound the table, flush when full
  4. several partials per key is fine
  5. measure records per key before writing it

basics

~20 s

Add an explicit first step that walks each piece of the input, accumulates per key in a bounded table, and emits one partial per key; then redistribute and combine the partials. It costs producer memory, a lookup per record, and a two-step shape a later reader must understand.

solid answer

~50 s

Write the aggregation as two steps instead of one. The first runs entirely inside each piece of the input with no movement: walk the records, keep a keyed accumulator table in the worker's memory, and at the end of the piece emit one partial result per key. The second performs the redistribution and combines the partials into final values. The traffic then falls from every record to at most one partial per key per piece, which is what a runtime-supplied fold would have achieved. The costs are real: a table whose size follows the distinct keys in the piece, a lookup for every record, and a pipeline that now has two steps where the intent had one. Bound the table and flush partials early when it grows, and skip the whole exercise when records-per-key is near one.

code

python · 16 lines
python
def fold_one_piece(records, max_entries=100_000):
    """Runs inside a single piece of the input. No movement happens here."""
    partial = {}
    for key, value in records:
        partial[key] = partial.get(key, 0) + value
        if len(partial) > max_entries:
            # table outgrew its budget: emit what we have and start again,
            # which yields more than one partial for some keys
            yield from partial.items()
            partial.clear()
    yield from partial.items()


def combine_partials(key, partials):
    """Runs after the redistribution, on the worker that owns this key."""
    return sum(partials)

go deeper

for a junior

The idea to remember is a two-step aggregation: combine what one machine already holds, then send those partial results to be combined again. It is the same answer with far less crossing the network.

for a middle

Describe the local step precisely — walk the piece, accumulate per key, emit one partial per key — and say why it adds no waiting line: it touches nothing outside its own piece and sits before an exchange that was happening anyway.

for a senior

Show the operational judgement: bound the accumulator table and flush early rather than risking a worker memory failure, measure records per key inside a piece before writing anything, and say which regime you are in because the flush point differs.

for a principal

The question worth owning is whether teams should be hand-writing this at all. A shared helper, or moving the step to a surface where the fold is inserted for you, removes a class of defect that otherwise reappears whenever key cardinality shifts.

## When you have to do it yourself Folding locally before the move — 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 — is sometimes applied for you and sometimes not. It is not applied when the step was written so that the per-key computation is defined only over a whole group, when the surface is a sequence of per-record functions executed close to literally, or in the oldest lineage of this class, where the local fold is something the author supplies rather than something a planner infers. In all of those cases the saving is still available; you just have to write it. ## The two-step shape 1. **A local step, inside each piece, with no movement.** A piece is one contiguous share of the input that a single worker processes on its own. Walk its records in order, maintain a keyed accumulator table in memory, and at the end emit one record per key: the key and its partial result. 2. **The redistribution.** Every worker sends each partial to whichever worker handles that key. 3. **A final combining step** on the collecting side, which combines the partials for a key into the finished value. The important property of step 1 is that it touches nothing outside the piece, so it adds no movement of its own and no stage barrier — no line across the job where downstream workers must wait for every upstream piece to finish. It is pure local work inserted before an exchange that was going to happen anyway. ## What it costs | cost | why it arises | how to contain it | |---|---|---| | producer memory | the accumulator table holds an entry per distinct key seen in the piece | bound the table; flush and clear when it exceeds the bound | | processor time | one table lookup and update per record | accept it; it is dwarfed by the network saving when the ratio is high | | more than one partial per key | an early flush emits several partials for the same key | harmless — the collecting side combines them anyway | | code that explains itself badly | the pipeline now has two steps where the intent has one | comment the first step as a volume reduction, not as logic | | nothing at all saved | records-per-key is near one, so each partial is one record | measure the ratio first and do not write the step | The memory row is the one that ends jobs. An unbounded accumulator table on a piece with many distinct keys is a worker-memory failure waiting to happen, and worker memory budgets and what happens when a working set exceeds them is a neighbouring subject. The defence is simple and belongs in the fold itself: cap the table by entry count and flush when the cap is hit. Emitting several partials per key costs a little more traffic and nothing in correctness. ## The regimes differ, so say which one you are in - **A finite job cut at a stage barrier.** The producing side is going to write one bucket per destination to local disk for the collecting side to gather later. A hand-written fold placed before that write shrinks both the local write and the later fetch, and the end of the piece is a natural, unambiguous moment to flush the table. - **A continuous job that hands each record over as it is produced.** There is no end of a piece and nothing is written between the halves, so folding by hand means deliberately holding records for a bounded count or a bounded time before emitting. That trades added delay for reduced volume, and it introduces per-key accumulators that persist across records — which is a subject of its own and should be designed deliberately rather than as a side effect of a volume optimisation. - **A continuous job run as a rapid succession of small finite jobs.** Each of those small jobs supplies its own natural flush point, so the hand fold looks much like the finite case, with the ratio computed over one small job's worth of records rather than over the whole stream — which weakens the saving when a key appears only once or twice per interval. ## Deciding whether to write it Before writing anything, estimate records per key: records entering the step divided by distinct grouping keys, and remember that the fold works within a piece, so what matters is the ratio *inside one piece*, not across the job. With a thousand pieces and a key that appears ten times in total, no piece sees enough of that key to fold anything. Rules of thumb that survive contact: - **High ratio, fold by hand; the movement shrinks by roughly that ratio.** - **Ratio near one, do not; you add a table and a lookup for nothing.** - **Ratio unknown, measure on a sample first** — the estimate is cheap and the rewrite is not. And if the per-key result genuinely cannot be built from partials — an exact median, an ordering of the group, cross-record comparisons — then no hand fold exists either, and the honest answer is that the full-volume movement is the work.

  • Why is capping the accumulator table safe, given that it emits several partials for one key?
    Because the collecting side combines whatever partials arrive for a key, and it makes no assumption that there is exactly one per producer. An early flush therefore changes only how many messages cross — more than the ideal, still far fewer than one per record. The cap is what keeps the fold from turning a network saving into a worker-memory failure.
  • How is this different from the fold a runtime would have inserted?
    Mechanically it is the same idea, applied at the same place. The differences are that you chose the memory bound and flush policy explicitly rather than inheriting one, that the step is visible in your source and must be maintained, and that nothing will withdraw it automatically if the key cardinality later rises and the fold stops paying.
  • Does a hand fold help in a continuous job that hands each record over as it is produced?
    Only by holding records back, since there is no end-of-piece at which to flush. You accumulate per key for a bounded count or interval and then emit, which buys volume at the price of delay and creates per-key accumulators that live across records. That is a deliberate design with its own state and expiry concerns, not a free local optimisation.

saying these in an interview costs you the question

  • Writes an unbounded accumulator table per piece and hopes
  • Thinks several partials for one key break the result
  • Adds the hand fold without checking records per key first
  • Assumes the local step introduces a wait for other workers
  • Believes a hand fold can produce an exact median per key