A nightly report program on one machine now takes seven hours; which checks come before moving it to a cluster?
answer
- measure before you distribute
- slow is not the same as too big
- one busy core, fifteen idle
- shrink the work before spreading it
- a squared cost survives more machines
basics
~20 sFind where the seven hours go before adding machines. A quadratic scan, repeated re-reads, and fields nobody uses stay expensive after distribution — and a run pinned to one core on a sixteen-core machine has headroom where it is.
solid answer
~40 sMeasure first: which resource is saturated, and for how much of the run. Three findings each point somewhere different. One busy core while fifteen sit idle means the headroom is inside this machine — parallelising within one address space costs no network and no coordination. Time dominated by parsing fields the report never uses, or by re-reading the same input several times, means the work itself should shrink. A cost that grows with the square of the input means the algorithm is the problem, and it stays the problem on a hundred machines. Distribution earns its place when the work divides into pieces that barely interact, every core on one saturated machine is already busy, and the deadline is still missed. Being slow is a symptom; it names no cause on its own.
go deeper
Recall that a long run has several possible causes and the fix depends on which one it is. Before anything else, find out whether the machine was busy or waiting, and whether it used one core or all of them.
Explain why some costs survive distribution and others do not: a quadratic cost keeps its growth curve on any hardware, wasted reading is simply spread around, while an idle-core finding means the capacity is already bought.
Demonstrate that you would produce the measurement before the proposal, name the three conditions under which distribution genuinely helps, and state what the move costs in operations, debugging and partial failure.
Consider what standard stops teams reaching for a cluster reflexively: who must show a profile before a platform change, what the default is for a new pipeline, and what the organisation pays when that default is wrong in either direction.
## Slow is a symptom, not a diagnosis "Seven hours" describes an outcome. It does not say whether the machine was working hard or waiting, whether one core was busy or all of them, whether the program read the input once or nine times, or whether the cost grows in proportion to the data or much faster than it. Each of those has a different cheapest fix, and only one of them is "more machines". The first thing to establish is **what was saturated and for what fraction of the run**: processor, storage read bandwidth, memory pressure, or a remote system the program was waiting on. ## The findings and what each one means | what you measure | what it means | cheapest fix | do more machines help? | |---|---|---|---| | one core busy, the rest idle | unexploited capacity on this machine | use the machine's cores within one process | yes, but you had a free order of magnitude first | | all cores busy, input read once | this machine is genuinely at its ceiling | a larger machine, or more machines | yes — this is the honest case | | most time parsing fields the report ignores | the program is doing work nobody asked for | take fewer fields and fewer rows | it spreads the waste rather than removing it | | the input read five times | repeated passes that could be one | restructure into one pass | each machine still re-reads its share | | cost grows with the square of the input | the algorithm, not the platform | fix the algorithm | no — the same growth curve, on more hardware | | long waits on a remote system | the program is not the bottleneck | batch or parallelise the calls | it multiplies pressure on that system | The interesting rows are the ones where distribution makes things worse. A run spread over many machines pays for its own coordination, and the waste is now paid many times over. ## Why the idle-core finding matters most A sixteen-core machine running one busy thread is using about six per cent of what it already has. Parallelising inside one address space has properties that spreading across machines does not: - no records cross a network; - there is no redistribution to pay for — no phase where each machine ships the records belonging to a given key to whichever machine is handling that key; - there is no partial failure: the run either finishes or it does not, rather than one machine of forty dying two-thirds of the way through; - the code and the way you debug it do not change. It is genuinely true that more machines could also parallelise this loop. The point is that you would be buying with money and operational surface something that is already sitting unused on the machine you own. ## Shrinking before spreading Two reductions repeatedly turn a seven-hour run into a one-hour run with no platform change: 1. **read less.** If the report uses four fields of eighty and one month of five years, the number of bytes the program touches falls by a factor that no amount of hardware will match. 2. **read once.** A program that makes five passes because five sections each want their own view of the input is doing five times the reading. Combining them into one pass that maintains all five accumulators removes four reads outright. These are not preliminaries to the real fix; very often they are the fix. ## When distribution is the honest answer All three conditions have to hold together: - the work **divides**: the input can be cut into pieces that are processed largely independently, with only a modest amount of cross-piece combination at the end; - the single machine is **already saturated** — every core busy, the input read once, nothing obviously wasteful left; - the deadline is **still missed**, or will be after foreseeable growth. If the computation is strictly sequential — each step's input is the previous step's output, over the whole data — then more machines add nothing at all, because no two steps can be in flight at once no matter how much hardware is available. ## What the move costs Adding machines is not a neutral act. The program now runs where you cannot attach to it; records that must meet each other cross the network; the run can fail partially and be retried partially; and there is a bill and an operational surface that a single machine does not have. How much is charged before the first answer **varies with how the machines are supplied** — a pool that is always up starts immediately but bills while idle, while a cluster raised for one run costs a start-up delay each time. None of that is a reason never to distribute. It is the reason to know which of the six rows in the table you are actually in before you do.
- The program is pinned at one busy core on a sixteen-core machine — what does that tell you?That roughly an order of magnitude of capacity is already paid for and unused. Parallelising inside one address space buys it with no network, no redistribution of records between machines, no partial failure and no change to how you debug the program. Spreading across machines could buy the same speed-up, but at a price you do not have to pay yet.
- Almost all seven hours are spent reading the input — what changes?The meter is now bytes read, so the moves are to read fewer of them: drop fields the report ignores, restrict the row range, and collapse repeated passes into one. If the program genuinely needs every byte within the window, aggregate read bandwidth becomes a real argument for more machines — but only if the input can be read by many readers in parallel.
- Why does a strictly sequential computation get no benefit from more machines?Because each step consumes the previous step's result, so only one step can be running at any moment. Adding machines increases how much could run at once, and the dependency chain means nothing does. The remedy is to find independence in the problem — per key, per partition, per time range — not to add hardware.
saying these in an interview costs you the question
- Proposes more machines without profiling where the time went.
- Believes distribution repairs an algorithm whose cost grows quadratically.
- Assumes every long-running program divides into independent pieces.
- Thinks a strictly sequential chain of steps runs faster on more machines.
- Ignores the unused fields and out-of-scope rows read every night.
- Treats a run with fifteen idle cores as evidence the machine is exhausted.