skip to content

A 40 GB event file must be totalled per customer on a 16 GB laptop — does that need a cluster?

level: juniorimportance: must knowfreq 74%

answer

  1. bytes at rest are the wrong meter
  2. ask what must be resident at once
  3. one entry per customer, not per record
  4. read in passes, discard as you go

basics

~20 s

Usually no. Bytes on disk are not the test; the working set is. Per-customer totals hold one running number per customer, so the records can be read in passes and discarded, and only the totals must fit in memory.

solid answer

~50 s

The question to answer is not how large the file is at rest but how much has to be resident at the same time — the working set. Totalling per customer keeps one running number per customer, so ten million customers is a few hundred megabytes of entries while the 40 GB streams through and each record is discarded once it has been added. If the answer needs only a few fields and a few days, discarding the rest as early as the program can shrinks the pass further. A cluster execution engine — a system you hand a whole program to, which splits that program's work across many machines, runs the pieces and puts the results back together — earns its place when the working set will not fit one address space, or when one saturated machine cannot finish inside the deadline. `40 GB` on disk establishes neither.

go deeper

for a junior

Recall the distinction that decides this: how big the file is versus how much has to be in memory at the same time. A running total per customer keeps one number per customer, and the records stream past and are thrown away.

for a middle

Explain the mechanics: the memory bill is distinct customers times bytes per entry, independent of input size, and a computation whose retained state grows with the records seen behaves completely differently from one whose state is fixed.

for a senior

Show that you would ask for the deadline and the growth before answering, and that you would cut unused fields and out-of-scope rows before proposing any change of platform. Name what a cluster charges before it answers.

for a principal

Frame the trade: the cheap answer here sets the default for every new pipeline. Decide what the organisation's standard test is for reaching past one machine, and what it costs to have people reach for a cluster reflexively.

## Three quantities that get called "too big" The phrase "the data is too big for one machine" collapses three different measurements, and only one of them is a hard ceiling on a single process: - **bytes at rest** — what the input occupies in storage. This is the number everybody quotes and the weakest signal of the three. - **bytes resident at once** — the **working set**: the memory the computation genuinely has to hold simultaneously, which depends on what is being computed, not on the size of the input. - **wall-clock time** — how long the run takes against the window it has to finish in. Bytes at rest matter only indirectly: they set how long a read takes and therefore feed the third quantity. They do not, by themselves, say anything about the second. A 40 GB file and 16 GB of memory is not a contradiction unless the computation actually needs 40 GB resident. ## Why a per-customer total needs almost nothing The loop is: 1. read the next record from the file; 2. pull out the customer identifier and the amount; 3. add the amount to that customer's running number in a table of totals; 4. discard the record and go back to step 1. Nothing in that loop ever holds two records at once. The memory bill is the table of totals: **number of distinct customers × bytes per entry**. Ten million customers at roughly fifty bytes an entry is around 500 MB — comfortable inside 16 GB, and completely independent of whether the file is 40 GB or 400 GB. The input has been read in a single pass and thrown away as it went. ## What has to be resident, by computation | computation over the file | what must be held at once | grows with | |---|---|---| | sum, count, minimum, maximum over all records | one number | nothing | | running total per customer | one entry per customer | distinct customers | | exact count of distinct values | every distinct value seen | distinct values | | exact median of a column | more than a fixed amount, in general | the input | | a globally sorted copy of the file | more than memory, in general | the input | | lookup against a second, smaller table | that smaller table | the smaller side | The first two rows are the common case in reporting work, which is why so many "we need a cluster" conversations are really about the first column of the table and not the second. ## Shrinking the work before spreading it Before reaching for more machines, two moves often remove the problem outright: - **take fewer fields.** If the report needs three fields out of sixty, discarding the rest as early as the program can cuts both the working set and the time spent parsing. - **take fewer records.** If the report covers one month, reading one month rather than five years changes the arithmetic by a factor nobody can match by adding machines. Both shrink the work. Distribution only spreads it, and spreading work you did not need to do is the most expensive way to not do it. ## What would genuinely force more than one machine - the working set exceeds the largest single machine available — for example, the entry table itself is hundreds of gigabytes; - the deadline cannot be met by one machine even with every core busy and the input read once; - the input cannot be read fast enough by one machine's storage or network within the window; - growth will make one of the above true soon enough that the migration has to start now. ## What the cluster charges before it answers anything A cluster is not free at the moment you choose it, and what it charges **varies by how the machines are supplied** — this is exactly the kind of claim that is true of one arrangement and false of the next: - a **standing pool of machines that is always up** answers with little start-up delay, but it bills while nothing is running; - a **cluster raised for this one run and torn down after** bills only for the run, but pays a start-up delay before the first record is read; - **capacity you never size yourself** shows you neither number directly. On top of that, once work is spread, records that have to meet each other cross the network — a redistribution, where each machine sends the records belonging to a given customer to whichever machine is accumulating that customer. One process on one machine never pays that, because no two records are ever separated by a network. There is also more to operate and more to fail partially. So the honest answer to the interview question is a question back: what has to be in memory at once, and what is the deadline? Until those are known, 40 GB is just a number about a disk.

  • Which per-customer computation would break the single-pass argument?
    Anything whose per-customer memory grows with the records seen rather than staying at one number. An exact count of distinct products per customer has to remember which products each customer already bought; an exact median per customer has to retain the values. Both turn a fixed-size entry into an entry that grows, and the table then grows with the data rather than with the customer count.
  • The totals themselves need 40 GB and the laptop has 16 GB — is a cluster forced now?
    Not immediately. Three single-machine moves come first: a larger machine, since machines with hundreds of gigabytes of memory are ordinary; splitting the pass by customer-identifier range so each of several sequential passes holds only its share of the entries; or keeping the entry table on the machine's own disk. Each trades time for memory, which is often the cheaper trade.
  • Does it change the answer if the file lives in shared storage rather than on the laptop?
    It changes the read time, not the working set. Where the bytes sit affects how long the pass takes and how reliable it is, so it can push on the deadline question. It says nothing about how much has to be in memory at once, which is the test that decides whether one process can do the job at all.

Counting how many books each library member has borrowed does not require emptying the shelves onto one table. You read the loan ledger line by line and keep one tally per member. The ledger can be far longer than the table will ever hold, because the ledger passes across the table and only the tallies stay.

saying these in an interview costs you the question

  • Says any file larger than memory requires more than one machine.
  • Assumes the whole file must be loaded before any aggregation can start.
  • Confuses bytes stored on disk with bytes needed simultaneously.
  • Reaches for more machines before asking what the result must hold.
  • Believes more machines always shorten a run, whatever the work is.