You double the number of pieces the input is cut into and the same single unit still runs out of room — what does that finding narrow the cause to?
answer
- more destinations, same key
- the rule keeps one key together
- near-zero spill before the failure
- regroup before raising memory everywhere
basics
~20 sIt rules out volume that the division rule can actually divide, and points at something the rule holds together: one group that must land in one place, or one record too large. The next lever is what the data is grouped by, not how finely it is cut.
solid answer
~50 sCutting the input into more pieces shrinks each unit of work only where the failing unit's contents are divisible by the rule that assigns them. For a step fed by a **redistribution** — a step that cannot be computed from the records one worker already holds, so every worker writes its output split by destination and every worker fetches its share — the destination is chosen from the key, so every record carrying one key lands together however many destinations exist. If the failure survives a doubling, it is almost never raw volume; it is one group the rule refuses to split, or one individual record that is indivisible at any count. That redirects the remedy: change what the data is grouped by, or reformulate the step so nothing has to be held whole, before you raise memory on every machine. Confirm it by reading the failing unit's own reported numbers against a median unit's.
go deeper
Recall that records sharing a key are deliberately sent to the same place, so cutting the input into more pieces cannot make one key's share any smaller.
Explain the discrimination the experiment performs: divisible volume shrinks with the count, an indivisible group does not, and the reported per-unit numbers tell you which happened.
Show that you confirm rather than assume — unit count actually changed, failing unit against median, spill near zero — and then rank regrouping and reformulation ahead of buying memory.
The angle is what you standardise: whether teams have a supported way to express a whole-group operation incrementally, so this diagnosis does not have to be rediscovered per job.
## What the doubling actually tested Cutting the input into more pieces is the cheapest experiment available on this failure, and its value is that it discriminates. A unit of work is one worker's share of one step, and the question a doubling answers is whether that share is set by *volume the rule can divide* or by *something the rule holds together*. - If the unit was simply carrying more bytes than its budget, halving the volume per unit either fixes the run or visibly moves the peak. That is the volume case. - If the failure lands on the same logical unit with the same contents, the division rule has not touched the thing that was too big. That is the indivisible case, and it is where the interesting reasoning lives. A necessary caveat on the experiment itself: a doubling only takes effect on the step you changed. If the failing step is fed by a redistribution whose destination count you did not change, or if the source's own layout fixes how a read step is divided, the experiment may not have run at all — check that the reported unit count actually changed before you believe its result. ## Why more destinations do not split a key A **redistribution** is a step that cannot be computed from the records one worker already holds, so every worker writes its output split by destination and every worker fetches its share. The destination for a record is derived from its key, so that all records sharing a key meet in one place, which is the entire point of the exercise — a grouping cannot be computed if a key's records are scattered. That guarantee is what defeats the doubling. Doubling destinations halves the *number of keys* per destination; it does not divide any single key. Two consequences follow: 1. A unit whose bulk is many ordinary keys gets smaller at every doubling. 2. A unit whose bulk is one key gets no smaller at any count, forever. Why one key carries such a share of the records, and the technique of spreading one key deliberately across several destinations, are a separate subject with their own costs; what belongs here is the inference — the failure surviving a finer division is the evidence that you are in case 2. ## Confirming it rather than believing it Read the run's reported numbers: whatever the engine reports per unit of work after a run — records in, bytes read, peak memory, bytes written to local scratch disk. Engines differ in how much they report and what they call it, but all report something per unit, and two comparisons settle this: | What you compare | Volume case | Indivisible case | |---|---|---| | Failing unit's records in, against the median unit's | Much larger, and it halves when you double | Similar, or large and unchanged by the doubling | | Bytes written to local scratch disk by that unit | Often substantial before it dies | Often near zero — nothing could be set aside | | Behaviour across attempts | Moves as the division changes | Same unit, same point, every time | The near-zero spill figure is the quiet tell. An operator that dies having written almost nothing out did not fail to keep up with a flood; it failed holding something it was not permitted to break apart. ## The remedy ranking, in order 1. **Change what the data is grouped by.** A finer natural key that carries the same meaning — the same subject plus a genuine dimension already in the record — makes every resident group smaller, and the parts are combined in a second step. This is the lever with the best ratio of effect to cost, and the one candidates reach for last. 2. **Reformulate so nothing must be held whole.** Replace possession of a group with a fold, a bounded top-N, or an approximation with a stated error. Spilling can rescue a fold; it can rescue nothing that must exist entire at one instant. 3. **Shrink the records before the grouping.** Drop fields the operation never reads. The group is then made of smaller records, which is a multiple on the same structure. 4. **Raise the memory every worker gets.** Last, deliberately, and with a note of why. It buys one multiple of headroom and charges it on every process in the cluster for the whole run. ## What this finding does not tell you - It does not tell you what the piece count *should* be — how the count is chosen in the first place is its own subject, balancing parallelism against per-unit overhead. - It does not tell you whether to retry the unit or the job; recovery granularity is a separate question with its own trade-offs. - It says nothing about whether the process was refused an allocation by the engine or killed outright by the platform at a ceiling the engine never saw. Those are two different failures wearing the same phrase, and which one you had changes who you talk to next. - Do not expect the engine to solve it for you. Some engines adjust a plan mid-run from what they measured, for some operators only, and a continuously running job largely not at all — so runtime replanning is a possibility to check for, never an assumption to rest on.
- The unit count you asked for did not change in the run's numbers. What now?Then the experiment never ran. The step's division may be fixed by the source's own layout, or the count you set applies to a different step than the one that failed. Confirm the reported unit count changed before drawing any conclusion from the outcome.
- The failing unit spilled a large amount before dying. Does that change your reading?Yes — substantial spilling says the operator did have a partial form and was using it, so the shortage is more likely volume plus an aggressive concurrency inside one worker than an indivisible group. Re-check how many units share that worker's budget.
saying these in an interview costs you the question
- Keeps raising the piece count after it stopped helping
- Believes finer division eventually splits a single key
- Assumes the engine detects and splits a heavy unit itself
- Jumps to a fleet-wide memory increase without reading the numbers
- Reports the doubling as a failure without checking it took effect