Who fixes the piece count of a distributed job, and when - at submission, mid-run, or only on restart?
answer
- no single regime across engines
- derived, author-set, revised, declared
- continuous width stands until restart
- restart moves per-key state across lanes
basics
~20 sIt varies by engine and by job kind. Some derive and freeze the count at submission; some let the author set it per redistribution point and adjust it between step groups while running; a continuous job's declared width usually stands until a restart, and some platforms never expose it at all.
solid answer
~40 sThere is no single answer, and saying so is the answer. In the older disk-to-disk lineage the split count is computed from the stored input before anything runs and is fixed for the run. In a modern finite job the author normally sets the count for each point where records are redistributed, and some engines will also revise it between step groups while the job runs. In a continuous job the author declares an operator width up front; because nothing about an endless input can be measured, that number usually stands until the job is restarted - and restarting at a different width means any per-key state has to be redistributed over the new lanes, which is not free. Some managed platforms derive it and expose nothing.
go deeper
Recall that the number of pieces is decided before or at the moment a job starts, and that it is a different number from how many machines or worker threads the cluster has.
Explain the regimes side by side and why they differ: an input that can be measured lets the count be derived, an endless input forces it to be declared, and a declared one usually stands until a restart.
Demonstrate that you plan for the regime you are on - correcting a finite job by resubmitting, but treating a width change on a continuous job as a scheduled operation with state redistribution and catch-up.
The angle is operability: which of your platform's workloads can be re-sized cheaply and which need a maintenance window, and whether that difference should steer new work toward one execution model.
## Why this question has four answers The **piece count** - how many slices the input is divided into, each processed by one **worker thread** from start to finish - is not owned by the same party in every engine. A candidate who answers with one regime is describing the engine they happen to know. The four regimes actually in the market: | regime | who sets it | when | cost of changing it | |---|---|---|---| | derived and frozen | the engine, from the stored bytes | before the run starts | resubmit the job | | author-set per redistribution | the author, in the program | at submission | edit and resubmit | | revised while running | the engine, between step groups | mid-run | none to the author | | declared width | the author, once | at submission | restart, with state moved | ## Derived and frozen The **two-phase disk-to-disk model** - the older execution model that materialises every intermediate result to disk between the two halves of a job - computes its splits from the stored input before a single record is read, and the number stands for the run. You influence it only through the input: how the files are laid out and how large they are. Nothing is adjustable once the run is in flight. ## Author-set at submission Most finite jobs on modern engines take a count from the author at each point where records must be redistributed. A **wide step** is one a worker cannot finish from the records it already holds, because it needs records currently sitting on every other worker; every such step ends one **step group** - the run of steps between two wide steps, scheduled and retried as one unit - and begins the next. The count applies to the pieces the next step group works on. Crucially, in several engines one number applies to *all* such points in the job unless the author overrides them individually, which is why one badly chosen value can be wrong in two directions inside the same run. ## Revised while running Some engines measure what a finished step group actually produced and change the count for the next one. That capability is **runtime replanning**, and it exists for some operators in some engines and largely not at all in continuous jobs. It is owned by Replanning at Runtime; the only thing to carry here is that it makes the initial number a starting point rather than a verdict, and only on the engines that have it. ## Declared width, standing until restart A **continuous job** is a run over an input with no end, so nothing about the input can be measured in advance and nothing can be derived from it. The author states a **declared operator width**: how many copies of each operator run. That number usually stands for the life of the job. Changing it means stopping and restarting, and that is where the real cost sits: 1. The job has to be stopped at a point whose position in the input is recorded, so it can resume without losing or repeating work. 2. Any state held per key has to be redistributed: the rule mapping a key to a lane depends on the width, so at a new width most keys belong somewhere else. 3. The job restarts behind the live input and has to catch up. Some engines make this a supported operation with tooling; others make it an outage. Either way it is a planned event, not a knob. ## Exposed to nobody A fourth case is growing: a platform that accepts work and derives the division itself, exposing no control at all. There your only levers are indirect - how much data the step actually has to handle, and the shape of what the previous job left on disk. ## What to say in an interview Three moves make this answer strong: - **Name the regimes rather than a setting.** *Fixed at submission here, per redistribution point there, restart-only for a continuous job* beats naming the one place your engine keeps the number. - **State the asymmetry.** A finite job can usually be corrected by resubmitting it; a continuous job cannot, so its width deserves more care up front. - **Say what you would still control** when you cannot set the number: the volume the step has to handle, and - by handing it to the right owner - the layout the previous job wrote, which is The Output Shape. ## The trap in the question The word *count* is doing three jobs, and only one of them is being set here. The **piece count** is the decision. The **worker-thread count** is a property of the cluster you were given. The **machine count** only decides how many lanes exist. An answer that drifts from the first to the third - *you change it by adding machines* - is the exact confusion this subject exists to correct: more lanes cannot help when the number of pieces is what is limiting the run.
- Why is changing a continuous job's declared width not simply a restart with a new number?Because the mapping from a key to a lane depends on the width, so at a new width most keys belong to a different lane. Whatever the job remembers per key has to be redistributed before the first record can be processed, and the job resumes behind the live input and must catch up. It is a planned operation, not a knob.
- If one number applies to every redistribution point in a finite job, what goes wrong?The job usually has steps of very different sizes. A value sized for the largest one shreds the smallest into tiny pieces whose fixed cost outweighs their work; a value sized for the smallest leaves the largest with pieces too big for one lane's memory. Where the engine allows per-point overrides, the large steps are the ones worth overriding.
- The platform exposes no control over the count at all. What is left?Indirect levers only: reduce how much data the step actually has to handle by filtering and aggregating earlier, and influence what the previous job leaves on disk, since file layout feeds the next run's division. Sizing the written output is owned by The Output Shape.
saying these in an interview costs you the question
- Every engine lets you change the count while the job is running
- A continuous job's width can be raised live with no consequence
- You change the piece count by adding machines to the cluster
- The engine always picks a good count, so the author never decides
- Restarting at a new width is free because the input is replayable