skip to content

Which computations over a 2 TB input refuse to run in bounded memory on one machine, and what still lets that machine finish them?

level: seniorimportance: should knowfreq 52%

answer

  1. fixed retention against growing retention
  2. sort and exact distinct resist one pass
  3. local disk is a second tier
  4. chunk, sort, write, merge the fronts
  5. count the passes against the window

basics

~20 s

Computations whose retained state grows with the input — a global sort, an exact distinct count at high cardinality, an exact median — resist bounded memory. One machine still finishes them by writing sorted chunks to local disk and merging them back.

solid answer

~50 s

Split computations by what they retain. Sums, counts, extremes, filters and per-key accumulators over a modest key count retain a fixed amount and run in one pass at any input size. A global sort, an exact distinct count over hundreds of millions of values, an exact median, and a match against a second input too large to hold all retain something that grows with the data — those are the ones that refuse bounded memory. A single machine still finishes them by using its own disk as a second tier: read as much as fits, sort or group it, write that chunk out, repeat, then merge the chunks in one streaming pass that holds only the front of each. The remaining ceilings are the disk's capacity, the number of passes multiplied by the read bandwidth against the deadline, and the largest machine you can obtain.

go deeper

for a junior

Recall that some computations keep only a fixed amount of memory no matter how much data arrives, while others must remember something that grows. Sums and counts are in the first group; sorting everything is in the second.

for a middle

Explain the chunk-sort-write-merge technique and why the merge holds only one record per chunk. Be able to say which computations it rescues and why an exact distinct count behaves like a sort rather than like a sum.

for a senior

Show the three ceilings you would actually measure — disk capacity, passes against the window, and the largest obtainable machine — and raise the requirements question of whether an approximate answer with a stated error is acceptable.

for a principal

Weigh the long single-machine run against the platform move: a six-hour run that fails at hour five has its own cost, and the decision is about risk and operational burden as much as about whether the technique fits.

## The dividing line: fixed retention against growing retention Every computation over a large input falls on one side of a single line: **does what it must keep in memory grow with the input, or not?** - **Fixed retention.** A sum keeps one number. A count keeps one number. A minimum keeps one value. A filter keeps nothing between records. These run in one pass over any input size, on any machine, and the input's size is irrelevant to memory. - **Retention bounded by something other than the input.** A running total per key keeps one entry per distinct key. This is fine as long as the key count is modest — it is bounded by the customers, the products, the regions, not by the records. - **Growing retention.** A globally sorted output cannot be produced without having seen everything, because the last record read may belong first. An exact distinct count must remember every distinct value already seen. An exact median must retain enough of the distribution to identify the middle. A match against a second input requires the other side to be reachable. The third group is the answer to the question. Note what is *not* on the list: a 2 TB input is not by itself a problem for the first two groups at all. ## The classes side by side | computation | what it must retain | one pass in bounded memory? | |---|---|---| | sum, count, minimum, maximum | one value | yes | | filter, reshape, derive fields | nothing between records | yes | | running totals per key | one entry per distinct key | yes, while keys are modest | | top ten by value | ten entries | yes | | exact distinct count | every distinct value | no | | exact median or exact percentile | enough of the distribution | no | | globally sorted output | effectively the whole input | no | | match against a second large input | the other side, or an index of it | no | ## The second tier: the machine's own disk A single machine is not limited to its memory. It has local storage, and the classic technique converts a memory problem into a time problem: 1. read records until memory is comfortably full; 2. sort or group that chunk in memory; 3. write the sorted chunk to local disk and free the memory; 4. repeat until the input is consumed; 5. open all the chunks at once and merge them, holding only the front record of each. The merge step is the important one: it produces a fully sorted stream while holding a number of records proportional to the **chunk count**, not to the input. A 2 TB sort on a 64 GB machine is roughly thirty-two chunks, and the merge holds thirty-two records plus buffers. The same trick handles the other members of the growing-retention group: an exact distinct count becomes a sort followed by a scan that counts value changes; an exact median becomes a sort followed by a positional read; a match against a large second input becomes sorting both sides on the matching field and walking them together. Alternatively, **partition by key range**: run the computation several times, each pass keeping only the keys whose identifier falls in one slice, so each pass's retention fits. ## Where the single machine genuinely stops Three ceilings, in the order they usually bite: - **local disk capacity.** The chunks plus the output have to fit. A sort needs room for roughly the input again. - **passes against the clock.** Each pass re-reads everything. Two passes over 2 TB at 1 GB per second is over an hour of pure reading before any computation is counted. Multiply by the passes the algorithm needs and compare against the window. - **the largest machine available.** This ceiling is finite, knowable in advance, and much higher than most people assume — but it exists, and once the retention exceeds it, no amount of patience helps. ## The third option: do not be exact Where the requirement allows it, approximate structures collapse the growing-retention cases back into fixed retention: an approximate distinct count and approximate percentiles both hold a fixed-size summary and report an answer with a stated error bound. This is a requirements conversation, not a technical trick — but it is the move that most often removes the ceiling entirely, and a candidate who never asks whether exactness is required has skipped it. ## What varies once you do reach for many machines Be careful about assuming every system in this class does the same thing. Many-machine systems that produce a globally sorted output run **the same chunk-sort-merge on many disks at once** — distribution buys parallel disks and parallel cores, not a fundamentally different algorithm. Systems designed around an input that never ends do not offer that operation in the same form at all. So the honest statement is that going wide multiplies the disks and the cores you can bring to a technique the single machine already had, and it is worth doing precisely when that multiplication is what you are short of.

  • Why is an exact distinct count harder than a sum over the same column?
    A sum retains one number regardless of how many records arrive. A distinct count must remember which values it has already seen, so its memory grows with the number of distinct values and nothing in the computation bounds that. The single-machine remedies are to sort and then count value changes, to partition by value range into several passes, or to accept an approximate answer with a stated error.
  • When does using local disk as a second tier stop being honest?
    When the chunks plus the output no longer fit the disk; when passes multiplied by the input size divided by read bandwidth exceeds the window the run has; or when the run is now long enough that a single failure late in it costs more than the migration you were avoiding. Each is a measurable threshold, so this can be checked rather than argued about.
  • Does a match against a second input always break bounded memory?
    No — it depends on the smaller side. If one input fits in memory, it is held as a lookup table and the large side streams past it in one pass. The bounded-memory case fails only when neither side fits, and then both sides are sorted on the matching field and walked together, which is the same chunk-and-merge technique applied twice.

saying these in an interview costs you the question

  • Says anything larger than memory must be spread across machines.
  • Believes sorting requires the whole input resident at once.
  • Treats an exact distinct count as no harder than a sum.
  • Forgets that every extra pass re-reads the entire input.
  • Rules out an approximate answer without asking whether exactness is required.
  • Assumes a machine's local disk is not usable as a second tier.