A candidate input is 400 MB compressed in shared storage - what must you know before copying it to every worker?
answer
- stored size is not memory size
- decoded and indexed, not compressed
- then multiply by the worker count
- estimates predate growth and filters
- leave headroom for the streaming side
basics
~20 sWhat it becomes in memory, not what it occupies on disk. Decoded into records and indexed for lookup, a compressed file commonly expands several-fold, and that expanded figure is what every worker must hold at once.
solid answer
~50 sFour different numbers get confused here. The 400 MB is stored, compressed and encoded in shared storage - durable storage every worker can reach. What matters is the decoded footprint: the same data as records in a worker's memory, plus the keyed lookup built over it, which is commonly several times the stored figure. How many times over depends on the encoding and on how the runtime represents a record - a compact binary layout the runtime manages is far smaller than ordinary language objects for the same row. Then multiply: that decoded figure is resident on every worker at once and sent to each of them. Finally, check the provenance of the estimate itself, because recorded statistics can predate the data's growth, can predate the filters you applied, or can be absent entirely for an input produced by an earlier step in the same job.
go deeper
Take away one rule: the size of a file on disk is not the space it needs in memory. Compression and encoding are undone when the bytes become records.
Be able to walk the chain - stored bytes, decoded records, the lookup over them, times the number of workers - and say which link each number belongs to.
Demonstrate suspicion of the estimate itself: where it came from, when it was captured, whether it predates the filters, and what happens when the input is produced by an upstream step nobody has run yet.
Ask what makes this decision decay safely. A margin that closes silently as data grows is an outage scheduled for later, so say who measures it and what the job does when the copy no longer fits.
## Four numbers, routinely swapped for one another The single most common mistake in judging whether an input can be copied to every worker is comparing the wrong pair of numbers. | Number | What it is | What it is for | |---|---|---| | Stored size | 400 MB, compressed and encoded in shared storage | How long it takes to read, and nothing else | | Decoded footprint | The same data as records in one worker's memory, plus the lookup built over it | Whether the copy fits at all | | Wire cost | The decoded or encoded form, delivered once per worker | What the copy costs the network | | The worker's budget | Everything that process must hold, of which the copy is one part | Whether it fits *beside the rest of the work* | Only the second tells you whether the copy is possible. The fourth - how one worker's memory is divided among the operators running in it - is a subject of its own; here it is enough to know the copy is not alone in that process. ## Why decoded is larger than stored, and by how much varies - **Compression is undone.** Whatever ratio the file achieved is repaid the moment the bytes become records. - **Encodings are expanded.** Values stored once and referenced repeatedly, or stored as small integers standing for longer values, become full values per record in memory. - **Record representation differs by engine.** Some runtimes hold rows in a compact binary layout they manage themselves; others hold ordinary language objects with per-object and per-field overhead. The same input can differ several-fold between the two, and no single expansion factor is true of the class. - **The lookup is not free.** Matching needs a keyed structure over the records - buckets, references, and slack so it stays fast - and that structure is charged on top of the records themselves. - **Only the columns you read count.** If the step reads three fields of forty, the decoded footprint is of three fields, which can turn an impossible copy into an easy one. So the honest answer to "how much will 400 MB become" is: measure it. If you must guess, guess high, and remember you are guessing about the runtime you are actually on. ## Where the estimate comes from, and how each source lies 1. **Statistics recorded alongside the data.** Cheap and immediate, and stale the moment the data grows. A job that has run nightly for a year may be relying on a figure captured when the input was a tenth of its present size. 2. **Metadata in the files about to be read.** Closer to the truth, but it describes stored bytes, not the decoded form, and it describes the whole input rather than what survives your filters. 3. **A measurement of a finished step.** The only figure that is actually true, and it is unavailable before that step runs. Some engines measure a completed step and revise the remaining plan on that basis; that revision is its own subject, and where it does not happen, an estimate made before anything ran is all there is. The third case is the one that bites: when the input to be copied is itself produced upstream in the same job - a filter, an aggregate, a set of distinct keys - nobody knows its size until it exists. ## Judging it in practice 1. Measure the decoded footprint once, on the runtime you actually use, for the columns the step actually reads. 2. Multiply by the number of workers to get the wire cost, and note the per-worker figure separately as memory held simultaneously. 3. Leave real headroom: a copy that consumes most of a worker's memory leaves nothing for the work streaming past it, and that work is why the copy exists. 4. Re-check on a schedule. This decision has a shelf life, because the input grows and the cluster is resized, and nothing warns you as the margin closes. ## The trap in one sentence A stored size is a property of a file; a copy is a property of memory in hundreds of places. "400 MB" is not small until you have said 400 MB of what, decoded into what, held by how many workers, alongside what else.
- The input to be copied is produced by an earlier step in the same job. How do you size it?Before that step runs, you cannot: any plan-time figure is a guess. Some engines measure a finished step and revise the rest of the plan, which is a separate subject; without that, measure a representative run first, or write the job so it still works when the copy turns out too large.
- Would sending the copy in its compressed form solve the memory problem?No. It reduces the bytes on the wire, but a worker cannot match against compressed bytes - it must decode them to answer key lookups. Transfer encoding changes the delivery cost, not the resident footprint, which is the number that decides whether the copy fits.
- Why can two engines disagree about whether the same input is copyable?Because they represent records differently. A runtime holding rows in a compact binary layout it manages may fit an input that a runtime holding ordinary language objects cannot, for the same data and the same worker memory. Any expansion factor you carry in your head belongs to one of them, not to both.
saying these in an interview costs you the question
- Compares the file's compressed size straight against a worker's memory.
- Assumes every runtime holds records in the same compact representation.
- Trusts recorded statistics that predate the data's growth.
- Forgets the keyed lookup's own overhead on top of the records.
- Sizes the copy against a whole worker's memory with no headroom.
- Estimates the whole input when the step reads only a few of its columns.