A source step tracks its own read position rather than anything per key; how is that state redistributed when the width changes?
answer
- not every entry has a key
- attached to the instance, not the key
- no bucket means no automatic move
- the author declares how it composes
basics
~20 sThat is per-worker state: attached to one parallel instance of a step, not to a key. With no bucket to move, it is redistributed by a rule the author picks — the collected entries split across new instances, or every instance given the whole set.
solid answer
~50 sA read position is **per-worker state**: state attached to one parallel instance of a step rather than to a key — typically a source's read position or a small lookup copy given to every worker — and redistributed by its own rule when the worker count changes. Because there is no key, there is no bucket, so none of the key-bound machinery applies: the runtime cannot decide where such a value belongs, since it does not know whether the value is a list of assignments to divide, a set to merge, or a copy every instance should hold. So the author declares the composition rule. The two common ones are *collect and split*, where the entries from all old instances are pooled and dealt out to the new ones, and *give everyone the whole set*, which is correct only when holding a duplicate is harmless. Choosing the wrong one silently duplicates or drops work.
go deeper
Know that some of what a job remembers belongs to a key and some belongs to a worker instance, and that a source's read position is the second kind rather than the first.
Explain why the second kind needs a declared rule: with no key there is no bucket, so nothing mechanical can decide where the value goes when the instance count changes. Name the split-the-pool and copy-to-everyone rules.
Show the failure you would look for in review: a value redistributed by the copy-to-everyone rule when it was really a work assignment, producing duplicate output no downstream step removes, and per-worker values that grow with data rather than with width.
Set the boundary for the estate: what a job may legitimately keep per instance, and the point at which a growing per-instance value must be re-expressed as key-bound state or moved outside the job entirely.
## Two kinds of retained state, two different rules Not everything a job remembers hangs off a key. **Key-bound state** is state a step may read or write only under the grouping key of the record it is currently handling, held by the worker that owns that key, so the read needs no coordination with any other worker. **Per-worker state** is the other kind: state attached to one parallel instance of a step rather than to a key — typically a source's read position or a small lookup copy given to every worker — and redistributed by its own rule when the worker count changes. | | Key-bound state | Per-worker state | |---|---|---| | Addressed by | The current record's grouping key | The parallel instance itself | | Needs upstream | A redistribution by key | Nothing | | Typical contents | Running totals, join sides, deduplication entries | A read position, a small reference copy | | Size driver | Distinct keys times per-key bytes | Number of instances times a small value | | On a width change | Whole buckets move to new owners | A composition rule the author declares | | Who decides where it goes | The bucketing, mechanically | The author, explicitly | The last row is the answer to the question. For key-bound state the runtime has a complete recipe: a key's bucket is fixed, so a width change is a re-deal of buckets. For per-worker state there is no such recipe, because nothing in the value says which instance should hold which part of it. ## Why per-worker state cannot be bucketed Consider three instances of a source each holding a read position for the input pieces they were assigned. Halve the width to one instance. What is that instance's read position? There is no answer derivable from the values alone — they must be combined, and how they combine depends entirely on what they mean: - If the value is a **list of assignments**, pooling all the lists and dealing them out again preserves the meaning exactly. - If the value is a **small reference copy** identical on every instance, every new instance should simply receive a copy. - If the value is a **counter of work this instance did**, neither rule is obviously right; you may want a sum, and the runtime has no way to know that. So the author supplies the rule. That is not a gap in the design; it is the honest consequence of state that is not keyed by anything the runtime understands. ## The rules an author can pick, and what each gets wrong 1. **Collect and split.** All old instances' entries are gathered into one pool and dealt out to the new instances. Correct for assignment lists and read positions, because every entry ends up owned exactly once. Watch the pool size: it is gathered in one place during the restart, so a per-worker value that grew large makes restarts slow. 2. **Give everyone the whole set.** Every new instance receives every entry. Correct when the value is a reference copy, or when acting on a duplicate is idempotent. Wrong for work assignments: two instances would both process the same piece, and the effect is a duplicate in the output that nothing later removes. 3. **Keep what you had, drop the rest.** Surviving instances keep their own entries and entries from departed instances are discarded. This is almost always a defect rather than a choice — the discarded read positions mean input that is silently never read or re-read from the beginning. ## Sizing and the trap Per-worker state is usually described as small, and usually is. The trap is that its size is a function of the width, not of the data: a value that is a few hundred bytes per instance is invisible at ten instances and awkward at a thousand, and under *collect and split* the whole pool is assembled during a restart. If a per-worker value grows with the data rather than with the width — a cache accumulating everything an instance has seen, for instance — it has become an unbounded retained set with none of the controls key-bound state has, and it should be re-expressed as key-bound state or moved out of the job. ## Where engines differ - Some engines expose per-worker state as a general facility any step may use; others restrict it in practice to sources and sinks, where the read position and the pending-commit list live. - The set of composition rules offered differs. *Collect and split* and *whole set to everyone* are the two that appear in one form or another almost everywhere, but the names, the defaults and whether you may write your own vary. - Under **repeated small finite jobs** — the runtime cuts an endless input into short finite chunks and runs a complete job over each one — the read position is commonly tracked by the coordinating process in the snapshot metadata rather than held by the step, so the redistribution question is answered by the framework instead of the author. - In **the two-phase disk-to-disk batch model**, which keeps nothing between runs, neither kind of long-lived state exists, and a run's input assignment is recomputed from scratch each time. State the rule as *the author declares how per-worker values compose*, and then say which model you are describing. A candidate who asserts one engine's default as the universal behaviour has made the most common mistake on this subject.
- Why can't the runtime bucket per-worker state the way it buckets key-bound state?Because there is no key to hash. The value describes the instance's own assignment or holds a copy meant for everyone, and nothing in it says whether a width change should split it, merge it, sum it or duplicate it. Only the author knows what the value means, so only the author can supply the composition rule.
- What goes wrong if every instance is handed the whole collected set after a width change?If the entries are work assignments, every instance now claims all of them and the same input is processed several times, producing duplicates that nothing downstream removes. The rule is only safe when the value is a reference copy or when acting on a duplicate is idempotent. It also makes the per-instance footprint grow with the width.
saying these in an interview costs you the question
- Thinks every retained value is bound to a grouping key
- Expects per-worker state to move with key buckets
- Hands every instance the whole list without checking for duplication
- Uses per-worker state to hold a per-entity running total
- Assumes per-worker state is always too small to matter