A program is handed to a cluster and reads no input for forty seconds. What is the cluster doing?
answer
- paid before any data is touched
- room, then processes, then the code
- listing costs per object, not per byte
- per job versus per machine
- amortised by long runs, repaid by short ones
basics
~20 sBefore the first read, a cluster must be granted capacity, launch a worker process on each machine, copy the program and its libraries to each, build a plan and list the input. None of that touches data.
solid answer
~50 sThe gap is real work, not network delay. The layer that owns the machines has to admit the job and reserve room for it, which can be a queue; a worker process has to start on each machine and report itself ready; the program's code and libraries have to be copied to every one of them; the coordinating process has to turn the program into steps and decide how the input is cut; and something has to list the input, which costs in proportion to the *number* of objects rather than their size. Admission, planning and listing are per job, so more machines do not divide them, while launch and code distribution are per machine and get slower as the cluster grows. Where workers are already up the first two shrink to almost nothing; where machines are acquired for this run, acquisition can dominate everything else.
go deeper
Be able to say that a cluster pays for capacity, process launch, code distribution, planning and input listing before it reads anything, and that a short run can therefore lose to one machine.
Explain which parts are per job and which per machine, why that makes the delay flat or rising as the cluster grows, and why a long-lived run amortises the whole list while a run repeated every few minutes does not.
Show that you measure it: time from hand-off to first read, compared against total wall clock over runs of different sizes, and act on the ratio by consolidating work rather than adding capacity.
Decide what this fixed cost means for a platform's defaults — the minimum run size worth submitting, whether capacity stays warm and who pays for it while it waits, and what the standing pipeline shape should be so the fixed part stays a small share of the bill.
## Start-up is work, not latency A **cluster execution engine** is 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. **Handing the program to the cluster** means packaging it together with the resources it is asking for and sending it off; from that moment you no longer choose which machine runs what. The gap between that moment and the first record being read is not wire delay — per-hop latency inside a cluster is measured in fractions of a millisecond. It is a sequence of real operations, each costing wall-clock time, and most of them happen on machines that are already billing. ## What happens in the gap 1. **Capacity is granted.** The layer that owns the machines — a pool of capacity that grants a job a share of itself and takes the share back when the job ends — must admit the job and reserve room. If the pool is busy this is a queue, and from the job's point of view the wait has no upper bound. 2. **Worker processes start.** A **worker process** is one process on one machine that runs pieces of the job and holds their data in its own memory for as long as it is alive. Each has to be launched and has to report itself ready before it is given work. 3. **The program travels.** Your code, the libraries it depends on and the settings it was given must reach every machine that will run a piece of it. A single process already has its code in memory; a cluster copies it once per machine. 4. **The plan is built.** The **coordinating process** — the one process that holds the job's plan, hands pieces of work to the others and collects what comes back — turns the program into a set of steps and decides how the input is cut into pieces. 5. **The input is enumerated.** Before anything is read, something must discover what is there: which objects exist, how large each is, how they divide. Against **the shared storage the cluster reads** — storage every machine can reach, which is what lets a piece of work run wherever there is room rather than where the bytes sit — this cost grows with the *number* of objects, not their total size. Half a million small objects can take longer to list than a dozen large ones take to read. 6. **The runtime warms.** Code paths that are compiled, cached or memory-mapped on first use are slower in the first pieces than in the thousandth, so early work is not representative of steady-state speed. | cost | one process on one machine | a cluster | |---|---|---| | getting capacity | already running | admission, possibly a queue | | starting processes | one, already started | one per machine | | making the code available | already in memory | copied to every machine | | planning | nothing worth naming | steps, plus a division of the input | | listing the input | one local listing | the same listing, usually remote | | warm-up | paid once | paid on every machine | ## More machines does not mean a shorter start This is where candidates go wrong. Admission, planning and input enumeration are **per job**: they take about the same time on four machines as on four hundred. Process launch and code distribution are **per machine**: they get slower, not faster, as the cluster grows, even though they happen concurrently. So the time before the first record is flat at best and rising at worst as capacity is added, while the compute after it falls. A run over a small input can therefore finish behind one process on one machine while the cluster is behaving perfectly — the fixed part simply dominated. ## What varies between engines and supply models - Where the workers are **already up** when the program arrives, the first two steps shrink to almost nothing, but the third is often still paid: the new program's code is not there yet. - Where machines are **acquired for this run**, acquisition can dominate everything else, and it is charged from the moment the machines exist, not from the first record. - Where capacity is **never shown to you**, there is still a start-up; it is inside somebody else's service, where you can neither watch nor shorten it. - A run that **stays up for weeks** pays the whole list once and amortises it into irrelevance. A runtime that instead collects whatever arrived during a fixed interval and then executes an ordinary finite run over just that slice, over and over, pays part of the list on every slice — one reason such an interval has a practical floor that has nothing to do with the data. ## What to say in an interview Name the components, say which are per job and which are per machine, and give the measurement: time from handing the program over to the first input byte read, against the run's total wall clock. If set-up is a third of the run, the work is being run in a shape it is too small for, and the lever is fewer and longer runs, not more machines.
- Why can a directory of very many small input objects make the delay worse than a much larger input in a few objects?Enumeration is per object: something must discover each object's existence, size and division before work is handed out, and that cost tracks the count rather than the bytes. Very many small objects also produce very many tiny units of work, each with its own hand-out and bookkeeping, so the fixed cost is paid again at a finer grain.
- Does a cluster that is already up remove this cost?It removes most of the first two parts — nothing has to be admitted from cold and no processes have to be launched — but the program's own code and libraries still have to reach each worker, the plan still has to be built, and the input still has to be listed. It shortens the delay; it does not delete it. The machines were also billing while they waited.
- How would you measure this rather than guess at it?Time the interval from handing the program over to the first input byte being read, and compare it with total wall clock across several runs of different sizes. The part that stays roughly constant as the input grows is the fixed cost; if it is a large fraction of a short run, consolidate work into fewer, longer runs.
Calling in a work crew. Before anyone lifts anything you have to get them released from other sites, get them through the gate, hand each one the drawings and the tools, and tell them who is doing what. Doubling the crew does not shorten the briefing, and for a job that takes ten minutes the briefing is the job.
saying these in an interview costs you the question
- Says the delay is network latency between the machines
- Thinks adding machines makes start-up shorter
- Assumes the program is already present on the workers
- Believes a cluster that is already up has no start-up at all
- Treats the time before the first record as free because nothing computes
- Thinks listing the input costs in proportion to its size