An input copied to every worker turns out far larger than estimated - where does the job fail, and why on many workers at once?
answer
- many workers fail together
- early in the step, before output
- not one unit, nearly all
- a retry repeats the same arithmetic
- shrink the copy or change the plan
basics
~20 sWherever the copy is assembled or held: on engines that gather it centrally, the gathering process can exhaust its memory first, and otherwise every worker fails at nearly the same point, because all of them are doing the identical thing.
solid answer
~40 sThe signature is correlated failure. Every worker - one process on one machine running some of the job's pieces - is holding the same oversized copy, so they run out of memory within moments of one another, early in the step and before much of the larger input has been read. On engines that pull the input back to the process coordinating the job before distributing it, that process can die first, which looks like a failure with no worker at fault. The distinguishing feature against an ordinary memory failure is the spread: one oversized piece of the larger input fails one unit of work on one worker, whereas a copy that does not fit fails almost all of them. Retrying is not a remedy, because the retry repeats exactly the same arithmetic.
go deeper
Remember the shape of the symptom: when many workers run out of memory at the same moment, suspect something every worker was asked to hold, not something one worker was unlucky with.
Explain why the failure is early and correlated - the copy must be resident before matching starts - and why a retry of the same units cannot change the outcome.
Separate it from a heavy piece and from a slow worker by spread and timing, then choose between shrinking the copy and abandoning the copy, and say what each costs.
Ask why the estimate was trusted without a guard. A plan that fails cluster-wide when an input grows needs either a measured margin or a fallback shape, and deciding which is a platform call, not a job-level one.
## What actually runs out Copying an input to every worker puts one whole copy in each worker's memory at the same moment. If the decoded copy is bigger than the estimate that justified the plan, nothing degrades gracefully: the memory is required all at once, before the step can match anything, and it is required in the same amount everywhere. ## The two places it fails | Failure site | When it happens | What it looks like | |---|---|---| | The process that gathers the copy | On engines that pull the input back to the process coordinating the job before distributing it | The job dies with no unit of work having failed; the workers were idle, waiting | | Every worker holding the copy | On any engine, once the copy is distributed | Dozens or hundreds of near-simultaneous memory failures, all early in the same step | Which of the two you see depends on the engine and on how it delivers the copy. Where each worker reads the input from shared storage itself, or where workers pass it onward to each other, there is no central gathering step to die - the failure is purely the second row. ## The signature: correlated, early, and repeatable Three things together identify this rather than any other memory failure: - **Correlated.** Many workers fail inside a short window. They are not failing at different points in different data; they are all holding the same object. - **Early.** The copy must be resident before the match begins, so the failure lands near the start of the step, with little of the larger input consumed and little output produced. - **Repeatable.** Rerun it and the same thing happens at the same place, because the input, the cluster width and the arithmetic are unchanged. Compare the two failures it is most often mistaken for: 1. **One oversized piece of the larger input.** One unit of work fails, on one worker, usually after a good deal of progress. A retry of that unit may even succeed on a less-loaded machine. That is a memory-pressure subject, not this one. 2. **A slow worker.** Nothing fails at all; the step simply waits. A straggler - one worker running slowly for reasons of its own - never produces a wave of simultaneous memory failures. ## Why a retry does not help A retry replaces a unit of work with an identical unit of work. Nothing about the copy's size, the number of workers holding it, or the memory each has, is different on the second attempt. Where a job is configured to retry a failed unit several times, the only effect is that it takes several times as long to fail - and on a continuous job, where the copy is rebuilt when the job starts, the same failure returns on every restart, which reads as a job that will not stay up rather than as a sizing mistake. ## What you actually do 1. **Copy less.** Reduce the input before it is replicated: apply the filter first, keep only the fields the match needs, and deduplicate where the match allows it. This is usually the cheapest fix by a wide margin, because every byte removed is removed once per worker. 2. **Stop copying.** If the input is no longer small, the match wants a strategy for two large inputs instead. That is a different subject, but recognising that the input has crossed the line is part of this one. 3. **Change the shape, not one machine.** Adding memory to one machine fixes nothing: the requirement is on every worker at once. Either every worker grows, or the copy shrinks, or the plan changes. 4. **Re-measure rather than re-guess.** The estimate that justified the copy was wrong; replacing it with another estimate from the same source repeats the mistake. ## The slow version of the same failure It does not always fail outright. A copy that fits, barely, leaves almost nothing for the records streaming past it, and the worker spends its time reclaiming memory instead of matching. The job then runs - far slower than it should, on every worker equally, with memory pressure high everywhere and no single culprit visible. Uniformity across the whole cluster is the clue: a problem in the data affects some workers, and a problem in the plan affects all of them. ## What changed since last month When a job that worked begins failing this way, the cause is nearly always one of three, and none of them is a code change: the copied input grew, the cluster was widened so more copies are held at once, or an upstream change made the input decode into more memory than before.
- How do you tell this apart from one oversized piece of the larger input failing?By the spread and the timing. An oversized piece fails a single unit of work on a single worker, usually well into the step; a copy that does not fit fails many workers within moments of each other, before much output exists. The first is about how the data is distributed, the second is arithmetic.
- The same job succeeded last month on the same cluster. What changed?Typically one of three things: the copied input grew, the cluster was widened so more copies are resident at once, or an upstream change made the input decode into a larger in-memory form. Any of them moves a job across the line with no code change at all.
- Would giving every worker more memory be a legitimate fix?It can be, but price it first: the extra memory is bought on every worker, not one, so it is the most expensive of the available fixes. Reducing what is copied removes bytes once per worker too, and costs nothing.
saying these in an interview costs you the question
- Treats it as a slow worker and waits for the step to catch up.
- Retries the job unchanged and expects a different outcome.
- Assumes only one unit of work failed, as with one oversized piece.
- Adds memory to one machine rather than to every worker.
- Blames the largest input, which in this shape never moved at all.
- Reads simultaneous failures as a cluster-wide hardware fault.