When the machine holding a piece's bytes is busy, should the scheduler wait for it or start elsewhere?
answer
- two costs in different currencies
- idle lane now against slow read throughout
- short pieces should not wait
- ordered rungs, each time-bounded
- vacuous when nothing is nearer
basics
~20 sIt depends on how long the piece will run. Waiting buys a cheap local read but burns an idle lane now; starting further away pays a remote read for the piece's whole life. Short pieces should never wait long.
solid answer
~50 sThe choice is between two costs that are not in the same currency. Waiting costs capacity that is idle right now, in a worker thread that could be doing something. Starting on a machine that does not hold the bytes costs a remote read for the entire life of that piece. The break-even is roughly: wait only if the expected wait is short against how long the piece will run, and only if remote reading is materially slower here. Schedulers that model this at all do not wait indefinitely — they hold a bounded moment for the machine that has the bytes, then widen the circle of acceptable machines, and eventually accept anywhere. Others express no preference at all and simply place on the first free lane, which is the right behaviour when no machine is nearer than another.
go deeper
Recall that there is a choice at all: start work now on a machine that must fetch the bytes, or hold on briefly for one that already has them. Both cost something.
Explain the break-even. The wait is worth it only when the piece runs long relative to the wait and the remote read is materially slower; say why short pieces should never wait.
Diagnose from symptoms. Idle lanes with a stretched wall clock means waiting too long; saturated network with collapsing per-lane read bandwidth means never waiting. Name the measurement that separates them.
Own the default across many jobs. A single cluster-wide wait policy is wrong for the spread of piece runtimes you actually run, and it is strictly harmful on any deployment where no machine is nearer than another.
## The choice, stated precisely A **piece of the input** is one slice that a single **worker thread** — one lane inside a worker process on one machine — reads and processes from start to finish. Suppose the bytes for that slice sit on the local disks of two machines, and both of those machines currently have every lane occupied. A lane is free on a third machine that does not hold the bytes. The scheduler has exactly two moves: - **Wait.** Keep the piece unassigned for a moment and hope a lane frees on a machine that holds the bytes. Cost: idle capacity now, on a lane that could have been working. - **Start now, further away.** Assign the piece to the free lane and let it fetch the bytes over the network. Cost: a remote read that lasts as long as the piece reads, not just at the start. Neither is free, and they are denominated differently — one in idle seconds across the cluster, one in network bytes and read latency for one piece. ## What actually decides it 1. **How long the piece will run.** A piece that finishes in two seconds cannot repay a five-second wait; a piece that runs ten minutes repays it easily. This is the dominant term. 2. **How large the local-versus-remote gap is here.** On a cluster where the network delivers close to local disk bandwidth, there is nothing to wait for. 3. **How likely the holding machine frees up soon.** If it is saturated for the next hour, waiting is a pure loss dressed as an optimisation. 4. **What the idle lane is worth.** On a busy shared cluster, an idle lane is someone else's queued work. On an idle cluster it costs almost nothing, so waiting is cheaper than it looks. ## Giving up in ordered steps A scheduler that models nearness does not make a binary choice; it walks an ordered set of rungs, from nearest to least near, dropping a rung each time a bounded wait expires: | rung | what it means | read cost | |---|---|---| | nearest | a lane on a machine whose own disk holds the bytes | local disk | | middle | a lane on a machine sharing a switch with one that holds the bytes | one hop, contended only within that group | | last | any lane anywhere in the cluster | across the shared fabric | The important property is that the rungs are **ordered and time-bounded**, not that any particular number of them exists. Some clusters have a meaningful middle rung because their network is hierarchical; some have a flat fabric where the middle rung is indistinguishable from the last one, and modelling it buys nothing. ## Where engines genuinely differ - **Some runtimes express an ordered preference and wait briefly at each rung. Others express none** and place every piece on the first free lane, because the deployment they target has no machine that is nearer. - In a **finite job** — a run over an input that ends, so the pieces are derived from measurable stored bytes — this decision is taken per piece, hundreds or thousands of times per run. In a **continuous job** — a run over an input with no end, where the author declares how many copies of each operator run — placement is settled once when the job starts and is not re-decided as records arrive, so there is no per-piece waiting to tune. - The older execution model that materialises every intermediate result to disk between the two halves of a job took this decision at a coarser grain, one round of pieces at a time, rather than continuously. ## The two failure modes - **Waiting too long.** Symptom: the cluster looks busy in the scheduler and idle in the machine metrics. Lanes sit unassigned while pieces queue for a small set of machines, and the run's wall-clock time is dominated by scheduling gaps rather than by work. This is especially sharp when the input's bytes are concentrated on few machines, so many pieces want the same few. - **Never waiting.** Symptom: every read is remote, the shared network saturates, and read bandwidth per lane collapses as parallelism rises — adding lanes makes the run slower, not faster. Both failures look like "the job is slow", and they are distinguished by one observation: whether lanes were idle or busy while the run dragged. ## What an interviewer is listening for That you treat it as an economic trade with a break-even, not as a setting. The strong answer names the quantity that sets the break-even — the piece's expected runtime against the expected wait, scaled by the read-speed gap — and says out loud that on a deployment where no machine holds the bytes, the whole decision is vacuous and waiting is strictly a bug.
- Why is one fixed wait applied to every piece a poor policy?Because the break-even moves with the piece. Pieces in the same run can differ by orders of magnitude in how long they take, and a wait that is well spent before a ten-minute piece is pure loss before a two-second one. A single number is right for one class of piece and wrong for the rest.
- What does a cluster look like when the wait is too long?Lanes sit idle while pieces queue, so machine-level utilisation is low while the run's wall clock stretches. The giveaway is that the elapsed time is not accounted for by work: sum the time pieces actually ran and compare it to the run's duration multiplied by available lanes.
- Does the same trade exist for a job whose input is remote storage?No. If no machine holds the bytes, every candidate lane is equally far, so there is nothing to wait for and the only sensible policy is first free lane. A preference configured in that environment can only cost idle time; it cannot buy a faster read.
saying these in an interview costs you the question
- Says waiting is free because the lane was idle anyway
- Claims a longer wait is always better since local reads are faster
- Thinks the scheduler can move the bytes to the free machine instead
- Treats every piece as equally worth waiting for regardless of runtime
- Assumes a scheduler waits indefinitely for the machine holding the bytes