skip to content

Adding Capacity Mid-Flight

Taking more workers or handing some back while the job is running: a finite job can usually absorb it at the next boundary, a continuous one must first move everything it remembers.

on this pageshow

questions

4

A job over input that ends gains extra worker processes halfway through its run. Where can that new capacity first do useful work?

level: juniorimportance: must knowfreq 58%

answer

  1. unstarted pieces only
  2. next wave, not the current piece
  3. boundary between rounds is safe
  4. gain capped by pieces remaining

basics

~20 s

On pieces of work that have not started yet, normally from the next round of work onward. A piece already running on another worker process is not cut in half and shared, and finished work is not redone to spread it more evenly.

solid answer

~40 s

New worker processes can only pick up work that is still free: pieces of the current wave that were planned but never started, and everything in later waves. A piece already in flight stays with the worker process running it, because the piece is the unit that gets handed out and it is not subdivided mid-run; finished work is not redone. That makes the natural landing point the boundary between rounds of work, when nothing is half-finished. A job over input that ends absorbs the change cheaply because the division of the work and the assignment of pieces are separate decisions, and the second is made late — the main thing that changes is how many places the remaining pieces have to run. The gain is capped by how many pieces are left.

go deeper

for a junior

Recall that extra worker processes take up work that has not started yet, and that a piece already running is never split between two processes.

for a middle

Explain why the boundary between rounds of work is the clean moment for the change: nothing is half-finished, so the next wave can simply be handed out to a larger set of processes.

for a senior

Show the judgment about when capacity is worth asking for: the gain is bounded by unfinished pieces and by the largest one, and a process granted near the end may never repay its startup.

for a principal

Frame the trade: elasticity on jobs with an end is nearly free to apply but pays back only on long runs with plenty of remaining pieces, so a platform policy should look at run length before it looks at load.

## The change this describes A **mid-run capacity change** is a running job acquiring more worker processes, or handing some back, while it is still running. A **worker process** is the process on a cluster machine that actually runs pieces of the job and owns the memory those pieces use; one machine can host several of them, so "four more machines" and "four more worker processes" are not the same count, and an answer should say which it is counting. Two neighbouring decisions are routinely confused with this one: - **Before the job is submitted** — how large each worker process is and how many are asked for. That is sizing, and it is settled before any piece runs. - **Underneath the job** — the platform's own control loop adding machines to the pool that many jobs draw from. Different actor, different lifecycle; the job sees it only as its requests being granted faster or slower. This question is the third thing: the job, already running, getting more places to run pieces. ## Where the new capacity can land A job over input that ends has a last round of work, and it gets there by cutting its input into pieces and running them in waves. A **round of work** is one wave of pieces that all run before the next wave can start, which gives the job a recurring moment when nothing is half-finished. New worker processes can take: - pieces of the current wave that have been planned but not yet started; - everything in later waves, from the boundary between rounds onward; - re-runs of pieces that failed or whose output was lost. They cannot take: - a piece already running elsewhere — the piece is the unit that gets handed out, and it is not cut in half in flight and shared; - work that already finished — a completed wave is not redone in order to spread it more evenly; - more pieces than remain — capacity beyond the number of unfinished pieces has nothing to do. | Candidate work | Can new capacity take it? | Why | |---|---|---| | Unstarted piece in the current wave | Yes, at once | Assignment happens whenever a place frees up | | Any piece of a later wave | Yes, from the next boundary | A wave is handed out after the previous one ends | | Piece already running | No | Pieces are not subdivided in flight | | Piece already finished | No | Redoing it adds work rather than removing it | ## Why a job with an end absorbs this cheaply Because two decisions are separate, and the second one is late. Cutting the input into pieces happens when the work is planned; deciding *which* worker process runs *which* piece happens as places free up. Adding worker processes changes the supply of places. It does not change the division of the work, it does not move anything that already exists, and it does not invalidate anything already computed. The job notices only that more pieces are in flight at once. That is the asymmetry with a job over input that never ends. Such a job keeps per-key state whose ownership follows the current worker count, so it has to relocate that state before it can run at a new width — a pause the job with an end never pays. ## What varies between engines - **Who asks.** Some runtimes request more worker processes automatically from a demand signal, such as pieces waiting with no free place, and hand them back after an idle period; others change width only when a person or a platform asks them to. - **What is granted.** The oldest two-phase disk-to-disk model in this class is granted places per unit of work rather than holding a long-lived set of processes, so "more capacity" arrives as more concurrent units rather than as new processes joining the job. - **Re-planning.** Some engines revise the division of later rounds from what a finished round measured. That is a separate mechanism from the capacity change; where both exist they interact, and neither implies the other. - **Startup.** A newly granted worker process must start, connect to the coordinating process — the single process that plans the pieces, hands them out and tracks what finished — and often read input for which it holds no local copy. On a short job that cost can exceed the work the new process does. ## The ceiling on what you gain 1. **How many pieces are left.** Forty new worker processes and three remaining pieces gives three busy processes and thirty-seven idle ones. 2. **How evenly the remaining pieces are sized.** One piece far larger than the rest sets a floor on elapsed time that no extra capacity lowers; diagnosing that unevenness is its own subject. 3. **When the capacity arrives.** Capacity granted near the end of a run pays its startup cost against very little remaining work.

  • What happens if the extra worker processes arrive when no unstarted pieces are left?
    They sit idle and the job finishes no sooner. Remaining elapsed time is set by the longest piece still running, and a job whose work was divided into few pieces simply cannot absorb the capacity. You pay for processes that do nothing, which is why capacity granted late in a run is often a loss.
  • Does a newly granted worker process start contributing immediately?
    No. It has to start, register with the coordinating process that hands out pieces, and usually read input it holds no local copy of. On a short job that lead time can be a large fraction of what remains, so the honest question is whether the job will still be running long enough to repay it.

saying these in an interview costs you the question

  • Thinks new worker processes take over a piece already running.
  • Expects finished work to be redone and spread across the new capacity.
  • Assumes doubling the worker count always halves the remaining time.
  • Believes the input is automatically re-divided whenever capacity changes.
  • Confuses the job asking for workers with the platform adding machines to the pool.
  • Forgets that a new worker process costs startup before it does anything.
open as a page

Why does changing the worker count of a job over endless input cost a pause, when the same change costs a job with an end almost nothing?

level: middleimportance: must knowfreq 55%

basics

~20 s

A job over endless input keeps per-key state, and which worker process owns a given key follows how many worker processes there are. Change that count and the retained state has to move to its new owners before any record can be processed correctly.

open as a page

A running job hands worker processes back to free capacity. Why is releasing one riskier than acquiring one?

level: seniorimportance: should knowfreq 45%

basics

~20 s

A worker process with no piece in flight can still be holding output that other worker processes have not fetched yet. Release it and that output goes with it, so the pieces that produced it may have to run again — often costing more than the release saved.

open as a page

Your platform can change the worker count of a job over endless input on demand. What makes an aggressive change policy cost more than it saves?

level: principalimportance: should knowfreq 32%

basics

~20 s

Every change stops the job while its retained state is relocated, and that cost follows how much the job holds rather than how much capacity moves. A policy that reacts faster than the pause is long spends its savings on pauses and provokes further changes.

open as a page