A stateful continuous job started at four workers must now run at forty; why can the key buckets fixed at its first start block that?
answer
- fixed at the first start
- a key never changes bucket
- whole buckets move, not keys
- the bucket count caps stateful width
basics
~20 sThe key space is cut into a fixed number of buckets at the job's first start, and a key never changes bucket. A width change therefore moves whole buckets between workers, and can never use more workers than there are buckets.
solid answer
~50 sKeys are not assigned to workers directly. They are assigned to **the fixed bucket count** — the number of key buckets fixed when the job first starts — and buckets are assigned to workers. A key's bucket is decided once by its hash and never changes, which is what lets a restart at a different width simply hand whole buckets to different workers instead of recomputing where every individual entry belongs. The consequence is a hard ceiling: with `N` buckets, at most `N` workers can own any state, and surplus workers own nothing. The count cannot be raised in place, because the stored entries were laid out under the old bucketing, so the routes out are to discard the retained set and rebuild it from a replayable source, or to rewrite the stored snapshot offline into the new layout. Engines differ in which of those they offer — and one of the three models in this family keeps no long-lived state at all, so none of it applies.
code
json · 17 lines{
"bucketCount": 8,
"firstStart": {
"workers": 2,
"bucketsOwned": { "w0": [0, 1, 2, 3], "w1": [4, 5, 6, 7] }
},
"restartAtWidth4": {
"workers": 4,
"bucketsOwned": { "w0": [0, 1], "w1": [2, 3], "w2": [4, 5], "w3": [6, 7] }
},
"restartAtWidth12": {
"workers": 12,
"outcome": "rejected by some runtimes, capped by others",
"reason": "8 buckets cannot be spread over 12 owners; 4 workers would own no key"
},
"invariant": "a key's bucket is decided once; only the bucket-to-worker map changes"
}go deeper
Know that a stateful job cannot always be made wider just by asking for more workers, and that something about how its stored values were laid out at the first start decides the limit.
Explain the indirection: keys hash to buckets, buckets are handed to workers, and a restart re-deals buckets. That is why entries move in whole groups and why the bucket count is a ceiling on stateful width.
Show that you treat the count as a day-one, irreversible decision: headroom over expected peak, written down with the deployment, and a stated recovery route — rebuild from source or rewrite the snapshot offline — if the guess turns out low.
Make it a standard rather than a per-job guess. Decide the organisation's default headroom and who reviews it, because the failure lands months later on someone who did not make the choice and cannot undo it cheaply.
## The bucket layer nobody notices until they need it There is an indirection between a key and a worker. A key hashes to a **bucket**, and a bucket is assigned to a worker. **The fixed bucket count** is the number of key buckets fixed when the job first starts: a key's bucket never changes, so a later parallelism change can only move whole buckets between workers, and can never exceed the bucket count. Why have the layer at all? Because without it, a width change would have to decide, per key, where its entry now belongs — a decision that would touch every entry in the retained set. With it, the width change is a reassignment of a small, fixed set of bucket identifiers, and the entries travel in whole groups that were already stored together. ## What is fixed and what is free | Quantity | Fixed when? | Changeable later? | |---|---|---| | The bucket count | At the job's first start | Not in place; only by rebuilding or rewriting the stored snapshot | | A key's bucket | By the key's hash, once | Never, for the life of the bucketing | | Buckets per worker | At each start | Yes, on every restart at a new width | | Number of workers owning state | At each start | Yes, up to the bucket count and no further | The first two rows are the irreversible ones, and they are the reason this comes up in senior interviews at all. Most tuning decisions in a processing job can be revisited next week. This one is made on the day the job first runs, usually without anybody noticing they made it. ## What a width change can and cannot do At four workers with, say, 128 buckets, each worker owns 32 buckets. Restarting the same job at forty workers is fine: the buckets are dealt out again, most workers get three, some get four, and each worker loads the entries for the buckets it has been given from the **durable snapshot** — the periodic consistent copy of the retained set written to storage outside the workers, which is what a restarted job reads. No key changed bucket; only the bucket-to-worker map changed. Restarting at 200 workers is not fine. With 128 buckets, at most 128 workers can own a bucket, and the rest own no keys and therefore no state. Runtimes differ in how they express that: some reject the start outright, some silently cap the effective stateful width, and some let the extra instances run as empty passengers that consume resources and do nothing. In every case the useful width stopped at the bucket count. ## Choosing the count on day one - **Set it from the width you might plausibly need, not the width you are starting with.** A count equal to today's worker count means the first scale-up is already blocked. - **Leave a multiple, not a rounding.** Headroom of several times the expected peak costs little and buys the one thing you cannot buy later. - **But do not set it absurdly high.** The bookkeeping is per bucket: each snapshot records more, smaller ranges, the metadata a restart reads grows, and very small buckets give a worker many tiny pieces instead of a few coherent ones. A large number is insurance, not a free lunch. - **Prefer a count that divides evenly by the widths you expect.** An uneven deal gives some workers one more bucket than others, which shows up as a lopsided load when buckets are similar in size. - **Write it down.** The number is a deployment fact, not a tuning knob, and the person who needs it is the one doing a capacity review a year later. ## Getting out of it when the guess was low 1. **Discard and rebuild.** If the source can be replayed far enough back and the retained set can be reconstructed from it, start a fresh job with a larger bucket count and let it rebuild. The cost is the replay and the period during which results are incomplete. 2. **Rewrite the snapshot offline.** Read the stored snapshot as an ordinary dataset outside the running job, redistribute its entries into the new bucket layout, and write a snapshot the next start can read. This requires the runtime to expose the stored form as readable data, and not all of them do. 3. **Run a second job.** Split the keyspace across two jobs by a coarse attribute. This is a workaround with real operational cost — two deployments, two snapshots, two sets of alerts — and it is chosen when the retained set cannot be rebuilt and cannot be rewritten. ## Where engines differ A record-at-a-time runtime typically exposes the bucket count as an explicit maximum-width number the author pins before the first start. In the **repeated small finite jobs** model — the runtime cuts an endless input into short finite chunks and runs a complete job over each one — the same ceiling appears as the redistribution width recorded alongside the state files at first start, and changing it is usually refused rather than negotiated. In **the two-phase disk-to-disk batch model**, which keeps nothing between runs, there is no long-lived retained set to relocate, and the next run is simply started at whatever width you like. Never state the ceiling as though every engine surfaced it the same way; state that it exists wherever long-lived key-bound state does.
- What are the options when the required width exceeds the bucket count?Discard the retained set and rebuild from a replayable source at a larger bucket count; rewrite the stored snapshot offline into the new layout, where the runtime lets you read it as data; or split the keyspace across two jobs. All three cost real downtime or real operational complexity, which is why the count is chosen with headroom on day one.
- Is a very large bucket count simply free insurance?No. Bookkeeping is per bucket, so every snapshot carries more ranges, the metadata a restart reads grows, and workers end up with many very small pieces rather than a few coherent ones. Pick a multiple of the width you plausibly expect rather than the largest number the runtime accepts.
- Does the bucket count also limit how many input pieces the job reads?No — those are different quantities. The bucket count governs where key-bound entries live and how many workers may own them. How the input is divided into units of parallelism is a separate decision, and a job can read far more input pieces than it has buckets.
A sorting office with a wall of numbered pigeonholes. Every street is assigned a pigeonhole on opening day and never moves to another one. Hiring more sorters just means handing whole pigeonholes to the new staff — but you can never usefully employ more sorters than there are pigeonholes, and widening the wall means taking down and re-sorting everything already on it.
saying these in an interview costs you the question
- Thinks width can be raised freely because keys rehash at restart
- Sets the bucket count equal to today's worker count
- Assumes a very large bucket count costs nothing
- Treats discarding the retained set as a routine way to change width
- Confuses the bucket count with the number of input pieces