skip to content

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%

answer

  1. state ownership follows the worker count
  2. stop, relocate, resume
  3. cannot process against half-moved state
  4. pause scales with retained volume
  5. not with how many workers were added

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.

solid answer

~40 s

A job over input that ends holds nothing across rounds that is tied to the worker count, so extra worker processes just take unstarted pieces. A job over endless input is different: it retains per-key state — running accumulators, buffered matches — and the mapping from key to worker process is a function of the current width. Change the width and many keys would land on a process that does not hold their accumulator, so the job has to stop taking records, relocate the retained state to its new owners, and only then resume. Processing against a half-moved accumulator would produce a wrong answer, which is why the move cannot overlap the work. The pause tracks how much is retained and how far it must travel, not how many worker processes were added.

go deeper

for a junior

Recall that a job over endless input remembers things between records, and that remembered material has to be moved when the number of worker processes changes.

for a middle

Explain the mechanism: key routing is computed from the worker count, so a changed count leaves accumulators on the wrong processes, and the job must stop, relocate and resume to stay correct.

for a senior

Demonstrate that the pause is sized by retained volume rather than by the size of the change, and plan around it — one large change, not many small ones, and a sense of what the resume costs.

for a principal

Take the position that width elasticity is a property a workload either earns or does not, and that jobs retaining a lot should be sized once for their peak instead of being adjusted continuously.

## Two jobs, one difference A **worker process** is the process on a cluster machine that runs pieces of the job and owns the memory they use. A **mid-run capacity change** is the job acquiring more of them, or handing some back, while it is still running. The interesting fact about this operation is that its cost is wildly asymmetric between the two job shapes in this class. - A job over **input that ends** reaches a last round of work. Between rounds it holds intermediate output, but nothing whose correctness depends on how many worker processes exist. A new process simply takes an unstarted piece. - A job over **endless input** never reaches a last round. To compute anything beyond a per-record transform it must retain state between records — a running total per customer, a buffer of unmatched events, a deduplication set. That retained state is spread across the worker processes the job currently holds. ## What actually has to move Records are routed to worker processes by their grouping key, and the routing is a function of the current width. Widen the job and a large share of keys are now routed somewhere other than where their accumulator sits. Before the first record can be processed at the new width, the retained state has to be relocated so that each key's accumulator is on the process that will now receive that key's records. The job cannot do that while it is processing: 1. **Quiesce.** Stop consuming input, and finish or hold what is in flight, so that no accumulator is being written while it is being copied. 2. **Relocate.** Move the retained state to its new owners across the network, or make it reachable from them. 3. **Resume.** Start consuming again at the new width, with input that arrived meanwhile now waiting. Processing during step 2 would mean updating an accumulator that is about to be overwritten by a copy of itself, or reading one that has not arrived — silently wrong numbers rather than an error. How engines arrange the key-to-process mapping so that less has to move is a genuinely separate subject, and so is what the retained set contains and where it is kept. ## The shape of the pause The pause is a data transfer, so its duration follows the **volume of retained state** and the distance it travels. It does not follow the size of the change: - going from eight worker processes to nine can cost nearly as much as doubling, because in both cases a large share of keys change owner; - a job retaining a few megabytes changes width in about the time it takes to start the new processes; - a job retaining hundreds of gigabytes can be out of service for minutes, during which input keeps arriving and has to be caught up afterwards. | | job over input that ends | job over endless input | |---|---|---| | What blocks the change | nothing; wait for the next boundary | relocating the retained state | | Typical cost | process startup | startup plus a stop-move-resume pause | | Cost scales with | nothing much | how much the job retains | | Correctness risk | none | processing against half-moved state | ## What varies between engines This is the point where a candidate who has only used one engine generalises wrongly: - **In place or on the way back up.** Some runtimes apply a new width to a running job, performing the relocation themselves. Others accept a new width only when the job is brought back up, so the change rides on a take-down-and-resume path — which is its own subject, and makes each change far more expensive. - **Continuous work as repeated small jobs.** One well-represented design runs endless input as a rapid succession of small jobs with an end. Such a job has a natural boundary every few seconds, and if it keeps nothing across those small jobs it can change width almost as cheaply as a job with an end; if it does keep per-key state across them, it still has to move it. - **Record-at-a-time with key-bound state.** Here the state is the defining asset, and the relocation is unavoidable and proportional to it. - **Width as a fixed property.** On several engines the width of a continuous job is, in practice, fixed for the life of the job, and "changing it" means a new run at a new width. ## What this means when you are asked it Say the asymmetry, then say the reason: the job with an end has late-bound assignment and nothing pinned to the worker count; the endless job has state whose ownership is computed from the worker count. Then say what varies, because an interviewer who runs a different engine from yours is listening for whether you know the pause is a property of retained state rather than of streaming as such.

  • Does the pause scale with how big the change is?
    Mostly no. It follows how much the job retains and how far that has to travel, so adding one worker process can cost nearly as much as doubling. That is the argument for making one large change rather than several small ones, and for treating width changes as rare events rather than a continuous adjustment.
  • Does a job over endless input that keeps no per-key state pay this?
    Barely. A per-record transform with no accumulators has nothing to relocate, so its width change is closer to the case of a job with an end: process startup, and re-reading whatever the new processes need. It still pays whatever the runtime's own mechanism costs if that runtime only accepts a new width on the way back up.
  • What arrives during the pause?
    Input keeps being produced whether or not the job is consuming it, so the pause leaves records waiting at the source and the job resumes behind where it was. Whether it then catches up is a separate question about throughput headroom, not about the width change itself.

Adding a room to a library does not help anyone until the books have been reshelved to match the new catalogue. Readers wait for the reshelving, and how long they wait depends on how many books there are, not on how big the new room is.

saying these in an interview costs you the question

  • Says the new worker processes just start taking keys as they arrive.
  • Thinks the pause is proportional to how many processes were added.
  • Assumes retained state can be rebuilt from records arriving after the change.
  • Claims every engine of this class changes width in place with no interruption.
  • Treats a job with an end and an endless job as paying the same cost.
  • Believes more worker processes always mean more throughput immediately.