An engineer sizes a worker's memory by multiplying input bytes by a guessed factor. What should replace that arithmetic?
answer
- input bytes are not memory bytes
- three sizes: at rest, packed, as objects
- measure, then size from the peak
- read per-unit numbers, not the job total
- peak has a date and a width on it
basics
~20 sA measured run. Read what the engine reported per unit of work - peak memory, bytes spilled, records in - against how many units shared the worker, then size from the observed peak plus room for bytes the engine never counted.
solid answer
~50 sThe arithmetic fails because the input's size is not the memory's size. The same records are compressed at rest, larger packed in whatever layout the engine holds them in, and larger again as ordinary language objects carrying the runtime's own per-object bookkeeping. What one worker actually needs depends on how much data each unit of work's slice carries, how many units run at once inside the process, which operators must hold a whole working set rather than streaming, and the bytes allocated outside the engine's accounting. None of that is recoverable from an input size. So run it - at reduced width, or over a representative slice - and read **the run's reported numbers**: whatever the engine publishes per unit of work after a run. Size from the observed peak with headroom, then measure again after any change to width, concurrency or the code.
go deeper
Recall that the size of an input on disk does not tell you how much memory it needs, and that the way to find out is to run the job and read what it reported.
Explain the three forms a record takes and the four things that set one worker's peak, then describe reading the per-unit numbers rather than the job total.
Show that you attach a width, a concurrency and a date to every peak you quote, and that you know a spilled run's peak was clipped by a region boundary.
Your angle is the standard: whether teams are expected to produce a measured peak before asking for capacity, and what the platform does when the measurement and the request disagree.
## Why input size does not predict memory The same data has at least three sizes, and they differ by an order of magnitude in both directions: - **compressed at rest** - what the storage bill and the file listing show; - **packed in the engine's own layout** - records held as a byte layout the engine allocates and interprets itself, where a field is read at a known offset; - **as host-runtime objects** - ordinary objects of whatever language the worker runs, each carrying the runtime's per-object bookkeeping, each reclaimed automatically, and collectively several times larger than the packed form. A factor multiplied onto the first number is a guess about which of the other two you land in, and engines differ on that: some hold packed bytes, others hold language objects, and the same program can sit in different forms at different points of the same job. That is before anything else on the list below. ## What actually sets one worker's peak - **How much data one slice carries.** A unit of work holds its own share, not the whole input. Cutting the input into more pieces lowers the peak per unit without changing the total. - **How many units of work run at once inside the worker.** Memory is owned per process and parallelism is per thread, so concurrency multiplies demand against one budget. - **Which operators must hold a whole working set.** A streaming projection holds one record at a time; an aggregation holds a table with one entry per distinct key; a join holds one side. The distinct-key count, not the record count, drives the aggregation. - **The bytes the engine never counted** - the language runtime, native and network buffers, and anything a user-supplied function allocates - which the machine holds even though the engine's report does not show them. ## The measurement 1. **Run it once, honestly.** Full width if you can afford it, otherwise reduced width or a representative slice - but a slice chosen so the distinct-key count and the heaviest key resemble production, because that is what the peak is made of. 2. **Read the per-unit numbers the engine reports**, not the job total. A job total hides the one unit that mattered; the distribution across units is the signal. 3. **Note peak memory alongside concurrency.** A peak recorded with two units per worker means something different from the same peak with eight. 4. **Read whether anything spilled.** Spilling - an operator writing part of its working set out to disk attached to the worker and reading it back to finish - tells you the working set exceeded its region. That is designed behaviour and often the right operating point, but it means the peak you observed was clipped by the region rather than being the operator's true appetite. 5. **Size from the observed peak plus headroom**, and leave the process room beneath the ceiling for the uncounted bytes as well as for the accounted budget. 6. **Re-measure after any change** to the number of pieces, the concurrency, the data volume or the user-supplied code. All four move the number. ## What the reported figures do and do not mean | figure | what it tells you | what it does not tell you | |---|---|---| | peak memory per unit | the high-water mark of what the engine accounted for | the process's real footprint, which includes uncounted bytes | | bytes spilled | the working set exceeded its region | how much room would have been enough to avoid it | | records in per unit | how evenly the work was divided | how many distinct keys, which is what a grouping table grows with | | job total memory | very little on its own | which single unit came closest to the edge | One caution on spill figures in particular: the same spill is commonly reported twice, once as the size it occupied in memory and once as the bytes actually written, and the two differ by however the engine encodes them. A large multiple between the pair is a fact about encoding, not a bug - so say which one you are quoting. ## What one measured run cannot tell you A measurement describes the data it saw. It does not describe a day with twice the volume, a key distribution that has shifted so one key now dominates, a width you have not run at, or a code change that added an allocation inside a user-supplied function. Treat the observed peak as an anchor with headroom and a date on it, not as a constant. Some engines also adjust their plan mid-run from what they measure, which moves memory pressure between runs of the identical program - that mechanism belongs to another subject, but its consequence here is that two runs are allowed to differ. ## The habit to demonstrate The interviewer is checking one thing: have you ever read a failed or a successful run's own numbers, or do you only know settings? A candidate who answers with a multiplier has told you they have not. A candidate who says *I would look at the per-unit peak and the spread across units, note how many units shared the worker, and leave room for what the engine does not count* has told you they have.
- Why does the measured peak move when you change how many units of work run at once in the worker?Because the budget belongs to the process and the units share it. More concurrent units means more operators holding working sets against the same wall, and each may also carry its own uncounted allocations. The per-unit peak can fall while the process footprint rises, so a peak figure is only meaningful next to the concurrency it was recorded at.
- What does a peak measured on a representative slice fail to capture?Anything the slice did not contain: a heavier day, a key distribution in which one key dominates, a width you have not run at, and any allocation a code change adds inside a user-supplied function. Sampling also tends to flatten the distinct-key count, which is exactly what an aggregation's in-memory table grows with.
- If a run spilled, is the observed peak the operator's real appetite?No. Spilling means the operator wrote part of its working set to disk attached to the worker because it could not hold it, so the peak you observed was clipped at the region boundary. The peak tells you the ceiling it hit, not how much it would have used if allowed - which is why the spill volume matters alongside it.
saying these in an interview costs you the question
- Multiplies input size by a factor and calls it a memory estimate
- Quotes a job total instead of the per-unit distribution
- Reads a peak without noting how many units shared the worker
- Assumes a peak measured once holds for every future run
- Treats compressed size and in-memory size as the same number
- Sizes only the accounted budget and leaves no uncounted room