skip to content

A 200 GB daily input that one large machine still handles has grown sixfold in eighteen months; how do you decide when to commit to a cluster, and what does deciding early or late cost?

level: principalimportance: should knowfreq 46%

answer

  1. extrapolate, do not measure today
  2. time to ceiling, not percent full
  3. two curves: records and distinct keys
  4. the largest machine is a knowable number
  5. set the trigger before the crisis

basics

~20 s

Decide on growth rate, not today's size: track projected months until the working set or the run time reaches what the largest available machine can do. Committing early pays a per-run floor cost; committing late migrates under a failing deadline.

solid answer

~50 s

Today's size answers nothing, because the decision is about the future. Convert the growth into one tracked number — **time to ceiling**: months until the working set, or the run time, reaches a stated fraction of what the largest machine you can actually obtain will do. Sixfold in eighteen months is roughly a doubling every seven months, so a machine with four times today's headroom is about fourteen months away from being the problem. Also check whether the working set grows at the same rate as the input; per-key state grows with distinct keys, which often grows far more slowly than records. Then set the trigger in advance rather than in the crisis, and keep the logic portable so the move is a change of runtime and not a rewrite. Early commitment pays a floor cost on every run and an operational surface you did not need yet; late commitment migrates while the nightly result is already late.

go deeper

for a junior

Recall that the question is about growth, not about today's number. A system that is comfortable now can be uncomfortable in a year, and the useful measurement is how quickly the input is getting bigger.

for a middle

Explain the extrapolation: a sixfold rise in eighteen months is roughly a doubling every seven months, so fourfold headroom is little over a year. Also separate record growth from distinct-key growth, since they bind different ceilings.

for a senior

Show the trigger you would set in advance and the numbers behind it — the largest obtainable machine's memory, cores, disk and read bandwidth — and describe the design constraints that keep a later migration cheap.

for a principal

Own the call and its reversibility: state what evidence flips the decision, who operates the result, what an early commitment costs every run, and why making the move cheap is worth more than predicting its date precisely.

## The decision is about a curve, not a point "One large machine still handles it" is a statement about today. The interesting question is when it stops being true, and that is set by the **growth rate**, which is the one input most teams never write down. Sixfold in eighteen months is a doubling roughly every seven months. Under that curve, the headroom people find reassuring evaporates quickly: | headroom today | months until exhausted at a seven-month doubling | |---|---| | 2x | about 7 | | 4x | about 14 | | 8x | about 21 | | 16x | about 28 | The table is the whole argument. "We are only at a quarter of the machine" sounds like years and is about fourteen months. ## Two curves, not one Before extrapolating, separate the quantities, because they frequently grow at different rates: - **records per day** — usually the number people quote, and the one that drives read time; - **the working set** — what must be held at once, which for keyed accumulation grows with the number of **distinct keys**. A business doubling its transactions has not necessarily doubled its customers. If the input doubles every seven months while distinct keys grow ten per cent a year, the run time is the thing approaching a ceiling and the memory is not — and the remedy is different. Extrapolating the wrong curve is how teams migrate a year early or a year late. ## The ceiling is finite and knowable The single-machine path ends somewhere specific, and you can look it up before you need it: 1. the largest memory you can actually obtain, from the suppliers you actually use, at a price the budget tolerates; 2. its core count, which bounds the parallel work inside one address space; 3. its local disk capacity, which bounds the second-tier technique of writing chunks out and merging them back; 4. its read bandwidth, which with the number of passes bounds the shortest possible run. These four are the real ceiling. Write them down once; they change slowly and they turn an argument into an arithmetic question. ## Track one number **Time to ceiling**: projected months until the binding quantity reaches a stated fraction — say seventy per cent — of the corresponding ceiling. Recompute it monthly from actual measurements rather than from the plan, and state the trigger in advance: - *begin the migration when time to ceiling falls below the time the migration takes, plus a margin.* That single sentence converts a recurring argument into a monitored threshold, and it is the deliverable this question is really asking for. ## What each mistake costs | | committing too early | committing too late | |---|---|---| | **run cost** | every run pays coordination and a fixed floor cost a single machine never charges | the last months run against a deadline that is already slipping | | **people** | everyone touching the pipeline must learn the distributed model | the migration happens under pressure, with the worst decisions made fastest | | **debugging** | work now runs on machines you cannot attach to, and failures are partial | unchanged until the move, then all of it at once | | **reversibility** | high, if the logic stayed portable | low, because nobody rewrites under a failing deadline | Note that what the cluster charges before it answers **varies with how the machines are supplied**: a pool that is always up starts immediately and bills while idle; a cluster raised for one run bills only for the run but pays a start-up delay before any answer; capacity you do not size yourself shows you neither. Do not cost the decision as if there were one arrangement. ## Keeping the decision reversible The cost of being wrong drops sharply if the code does not assume one address space. Practical constraints that keep the door open in both directions: - express the transformation as steps over key-partitioned data, so the same logic runs over one slice or over four hundred; - avoid depending on random access across the whole input, or on one global in-memory dictionary that only exists because everything happened to be in one process; - keep the input in a form that both a single process and a many-machine run can read; - keep the deadline, the working set and the growth rate as measured numbers on a dashboard, not as folklore. ## What would change the answer A principal-level answer names the conditions under which it flips early: the growth turns out to be a step change from a new source rather than a curve; a legal or contractual deadline shortens the window; the team that would operate the cluster does not exist yet, which pushes the trigger later and makes portability more valuable; or a requirements conversation reveals the daily job can work on a slice rather than everything, which removes the ceiling for another two doublings. The decision is not "cluster or not" — it is "what evidence moves us, and have we made moving cheap".

  • Which single number would you put on a dashboard to make this decision automatic?
    Time to ceiling: projected months until the binding quantity — working set or run time — reaches a stated fraction of what the largest obtainable machine can do, recomputed monthly from measurements. Pair it with one stated trigger: start the migration when time to ceiling drops below the migration's own duration plus a margin.
  • What makes a later migration cheap rather than painful?
    Logic expressed as steps over key-partitioned data, with no dependence on random access across the whole input or on a global in-memory structure that exists only because everything was in one process. Input kept in a form both a single process and a many-machine run can read. The move is then a change of runtime rather than a rewrite of the transformation.
  • The input doubles every seven months but distinct customers grow ten per cent a year — what follows?
    The two ceilings arrive at very different times. Run time and read volume are on the steep curve and will bind first, while the per-customer accumulator table is essentially flat and will not. The right response is a faster read path or more cores on one machine, and the memory-driven argument for more machines does not apply here at all.

saying these in an interview costs you the question

  • Decides on today's size and re-decides at every crisis.
  • Assumes a larger machine will always be available at any size.
  • Migrates at generous headroom with no growth measurement.
  • Assumes input growth and working-set growth are the same rate.
  • Treats the migration as a weekend of work with no rewrite.
  • Prices the move as if all machine supply arrangements charged alike.