skip to content

A pass reads a large file a piece at a time yet still exhausts memory totalling amounts per customer, and halving the piece size does not help. What is the peak footprint proportional to?

level: seniorimportance: should knowfreq 58%

answer

  1. two terms, one is not the piece
  2. the loop bounds the piece only
  3. one entry per distinct key
  4. cardinality times bytes per entry

basics

~20 s

Live bytes at the peak are roughly one piece plus the fold's accumulator. Halving the piece shrinks only the first term. Here the accumulator holds one entry per distinct customer, so it is proportional to the key's cardinality, which the piece size never touches.

solid answer

~50 s

Reading a piece at a time bounds **the piece**, not the computation. Peak live bytes are about `one resident piece + the fold's accumulator + whatever the step allocates transiently`, and only the first of those responds to the piece size. A total per customer keys the accumulator by the data itself, so its size is `distinct customers x (bytes of the key + bytes of the per-key state + the per-entry overhead of whatever lookup structure holds it)`. Note that the per-key state is tiny - a number, or a total and a count - the problem is how many of them there are. With a near-unique key the accumulator approaches the size of doing the whole thing at once, and the loop has bought nothing. The only levers on that second term are fewer distinct keys and fewer bytes per entry.

go deeper

for a junior

Remember the two terms: what one piece occupies while it is loaded, and what the fold carries between pieces. Only the first of those gets smaller when you shrink the piece.

for a middle

Be able to say what the accumulator is proportional to - one entry per distinct key, each holding the key's own bytes plus its state plus per-entry overhead - and why the piece size appears nowhere in that product.

for a senior

Do the arithmetic before touching anything: distinct keys times bytes per entry, against the cost of one resident piece. Then say which of the two terms exhausted the process rather than re-tuning the loop.

for a principal

The standing rule to own is which keys a fold is permitted to accumulate on. A key whose cardinality tracks the row count commits the team to a computation whose state is the dataset, whatever the loop around it looks like.

## Two terms, and the loop only touches one When a bounded input is processed a piece at a time in one process, the live bytes at the worst moment are approximately: `one resident piece + the fold's accumulator + the transient allocations of the current step` The piece size is a dial on the **first** term only. That is the whole content of the surprise: a run that dies with the piece set to a million rows dies again at a hundred thousand, and again at ten thousand, because the term that was too big was never the piece. Worse, shrinking the piece usually makes the run slower - more iterations, more per-piece fixed cost - so the engineer gets a slower failure and reads it as progress. The habit worth building is to say out loud what the accumulator is proportional to **before** running anything. ## What the accumulator is proportional to | fold | carried between pieces | grows with | |---|---|---| | a total over one column | one number | nothing | | a mean over one column | a total and a count | nothing | | a running smallest and largest | two numbers | nothing | | a total per customer | one entry per customer | the key's **cardinality** | The first three are fixed in size, and for those the piece loop really does bound the run. The fourth is keyed by values found in the data, and its size is: `number of distinct keys x (bytes of the key itself + bytes of the per-key state + per-entry overhead of the lookup structure)` Every part of that second factor deserves a number rather than a shrug: - **The key's own bytes.** A customer identifier of twenty-four characters is stored once per distinct key, and if each is held as an individual value rather than packed into one block, it carries a per-value header too. - **The per-key state.** Usually small - one number for a total, two for a mean. Assume it is small; that is not where the problem is. - **The lookup structure's overhead.** A hash-based structure keeps spare capacity and per-entry bookkeeping. This factor is a property of the design you are using, not a universal constant, and it can be several times the payload. Work it: eight million distinct customers at, say, forty bytes of key plus sixteen of state plus forty of per-entry overhead is around 768 MB of accumulator - against a piece of a hundred thousand rows that might cost tens of megabytes. The arithmetic says immediately which dial matters, and it says it before a single row is read. ## The worst case, stated plainly If the grouping key is close to unique - an event identifier, a request identifier, a session - then the number of distinct keys approaches the number of rows, and the accumulator approaches the size of doing the entire computation in memory at once. At that point the piece loop bounds the piece and bounds nothing else: the program has all the cost of the loop and none of its benefit. The converse is the reassuring case. A total per store across five thousand stores is a fold whose accumulator is a fixed few hundred kilobytes no matter whether the input is ten gigabytes or ten terabytes, and there the loop does exactly what it promised. ## The levers that actually exist 1. **Fewer distinct keys.** Coarsen the key, or restrict the pass to the keys you will actually report on, and the dominant factor shrinks in direct proportion. 2. **Fewer bytes per entry.** A narrower key representation and a more compact per-key state both attack the same factor, though by a constant rather than by an order of magnitude. 3. **Not the piece size.** It appears in the other term and nowhere in this one. If the key space is irreducibly near-unique and both levers are exhausted, the shape of the computation is what has to change, not its piece size. ## Combining two partial keyed results A pass that produces partial results and merges them at the end has one more moment to account for: during the combine, both inputs and the growing output are live at once, so the peak sits above the size of the final answer rather than at it. And how the two sides are paired depends on the design. Where the partials are **label-carrying** objects, the combine aligns on the labels and produces the union of both sides' keys, with an absence wherever one side never saw that key - and adding a number to an absence has to be given a meaning, or the totals for keys seen in only some pieces come out wrong or vanish. Where the partials are **positional** containers, values pair by position and the two sides must already be in the same order, which for a keyed fold means you are responsible for the ordering yourself. Decide which model you are in before you write the combine. ## Where designs differ How the per-key state is physically held varies widely: individual values in a general-purpose lookup structure, against one compact typed block per accumulator field addressed by a key index. The bytes per entry differ between those by a large factor, which moves the point at which a given input dies - but neither changes the proportionality. The accumulator is still one entry per distinct key, and that is still the term the loop does not bound.

  • The same pass survives on one input and dies on another of identical size. What differs between them?
    The number of distinct keys, not the bytes. A ten-gigabyte file of transactions across five thousand stores and a ten-gigabyte file keyed by device identifier are the same input size and two completely different accumulators. Size the run from the key's cardinality; the file's byte count tells you about the piece term and nothing about the other one.
  • Two partial per-key results are combined at the end of the pass. Where does the peak sit?
    Above the final answer, because both inputs and the growing output are live during the combine. There is a correctness point in the same moment: if the partials carry labels, the combine aligns on them and produces the union with absences where one side saw nothing, so what an absent partial means has to be settled before the addition.
  • Which changes actually shrink the second term, and which only appear to?
    Real: fewer distinct keys - coarsen the key, or keep only the keys you will report - and fewer bytes per entry through a narrower key and a leaner per-key state. Only apparent: a smaller piece, more iterations, or a faster loop. Those move the first term or the clock, and the accumulator is untouched by all three.

saying these in an interview costs you the question

  • Believes processing in pieces caps memory whatever the computation is
  • Halves the piece size again when the accumulator is what grew
  • Sizes the run from the file's bytes rather than the key's cardinality
  • Forgets the key's own bytes are stored once per distinct key
  • Assumes two partial keyed results always add position by position
  • Treats per-entry overhead of a lookup structure as negligible