A step attaches a 40 KB document to every record just before a grouping redistributes 200 million of them. Why can that movement dwarf the compute on both sides of it, and what do you change?
answer
- width multiplies everything after it
- arithmetic before tuning
- move the widening later
- drop unread fields earlier
basics
~20 sWidth multiplies every byte of the movement after it: 200 million records at 40 KB is about 8 TB on the wire, against 12.8 GB for the grouped fields alone. The repair is to widen after the movement, not before.
solid answer
~50 sDo the arithmetic before touching anything. Two hundred million records carrying roughly 40 KB each is about **8 TB** that has to cross the network, plus the same volume written to and read back from local disk wherever the producing side materialises its output, plus the processor time to encode and decode all of it. The grouping itself is a fold over a handful of fields — nanoseconds of work per record, so at most a few minutes of real computation. The movement is hours. The job is not slow at computing; it is slow at carrying bytes it need not have carried. The repair is to change what crosses, not how fast it crosses: attach the document after the redistribution if the grouping never reads it, carry only the fields downstream actually consumes, or move an identifier and re-associate later, paying a lookup instead of the bytes.
code
python · 13 linesrecords = 200_000_000
key_and_measures = 64 # bytes per record the grouping actually reads
attached_document = 40 * 1024 # bytes added by the widening step
def bytes_moved(width_bytes):
return records * width_bytes
wide = bytes_moved(key_and_measures + attached_document)
narrow = bytes_moved(key_and_measures)
print(wide / 1e12, "TB if the widening happens before the movement")
print(narrow / 1e9, "GB if it happens after")
print(wide / narrow, "times more bytes, for the same record count")go deeper
Recall that widening a record multiplies every byte of the movement that comes after it, and that a record count hides this completely.
Do the arithmetic aloud — encoded width times record count — and explain why placing the widening after the exchange, or the projection before it, changes the bytes rather than merely the shape of the plan.
Show the diagnosis and not just the fix: say how you would establish that the transfer rather than the computation is the job's time, and what you would verify before assuming anything was pruned for you.
Argue about where the widening belongs in the platform rather than in one job. A payload attached upstream is paid for by every consumer downstream of it, and moving an identifier instead trades bytes for a lookup that someone then owns.
## Do the arithmetic first A redistribution — every worker sending each record it holds to whichever worker will handle that record's key, so a grouping can see all of a key's records in one place — is costed as **encoded width times record count**. The widening step makes that first factor enormous: - key and grouped measures alone: roughly 64 bytes per record, so about **12.8 GB** moves; - with the attached document: roughly 41,024 bytes per record, so about **8.2 TB** moves. Same record count, a factor of about 641 in bytes. Where the producing side materialises its output, that 8.2 TB is also written to local disk and read back off it, and every byte of it is encoded once and decoded once. Nothing about the record count hints at any of this, which is why counting rows hides the defect completely. The compute on either side is the other half of the comparison. A fold over a few numeric fields is on the order of nanoseconds per record; 200 million of them is seconds to a couple of minutes of processor work spread across the cluster. When the movement is hours and the computation is minutes, **the job's duration is a transfer, not a calculation**, and every instinct aimed at the calculation will fail. ## How you recognise it - Compute bytes per unit of useful work. If the step consumes a handful of fields but the records carrying them are tens of kilobytes wide, the ratio itself is the finding. - Look at what changed the record between the read and the exchange. A widening is usually one innocuous line: a document attached, a nested structure exploded into repeated rows, an image or blob carried along "because it is needed at the end". - Watch the regime. Where the job is cut at stage barriers — a line across the job at which no downstream worker computes until every upstream piece of the producing step has finished — the movement appears as a long, visible phase between two short ones. Where records are handed downstream as they are produced, there is no phase to look at: the same problem shows up as the job's sustained throughput sitting below its input rate, which is a different subject to diagnose from there. ## The repairs, in order of bytes removed 1. **Move the widening to the other side of the movement.** If the grouping neither keys on nor reads the document, nothing requires it to travel: redistribute the narrow records, do the grouping, and attach the document to the far smaller result. This is the change that removes the factor of 641 rather than shaving it. 2. **Drop the fields nobody downstream reads, before the move and not after.** Only the bytes removed *before* the exchange are bytes the exchange never pays for. Whether this happens automatically depends entirely on the surface you wrote: a declarative pipeline may push a projection down for you, a program built out of opaque per-record functions gives the runtime nothing it can see through, and the oldest execution model of this class rewrites nothing at all. Never assume it was done; the projection you write yourself always crosses correctly. 3. **Carry an identifier instead of the payload.** Move a reference across the exchange and re-associate the document afterwards, trading tens of kilobytes of network for one lookup per surviving record — worth it exactly when the grouping collapses many records into few. 4. **Compress what must still cross.** A real lever, but a second-order one here: it charges processor time on both sides and cannot recover an order of magnitude if the documents are already compact. When many records share a grouping key, folding them together on the machine that produced them is the other major lever on a movement's size; it is a subject in its own right and is not what this question turns on. ## When the widening cannot simply be moved The document has to travel when the movement is keyed on something inside it, when a downstream filter reads its contents before the result narrows, or when re-fetching it later would cost more than carrying it. In that case the honest answer is not a trick but a decision: accept the movement, size the job for the bytes rather than the record count, and make sure it happens **once** — an expensive exchange whose output is reused by three downstream steps is a different economic proposition from the same exchange performed three times. ## What this is not It is not a tuning problem. Adding workers to a job whose time is bytes on the wire buys a little parallel bandwidth and nothing else; more memory does not reduce what has to cross; and raising the number of pieces the input is cut into changes the granularity, not the volume. The number that has to change is the width of the record at the moment it crosses.
- How would you tell that the movement, not the grouping, is where the job's time goes?Compare bytes against work. Estimate encoded width times the records crossing the exchange, then ask how much processor work each byte earns: a fold over a few fields is nanoseconds per record, so terabytes on the wire cannot be explained by it. Where the job is cut at barriers the movement is a visible phase between two short ones; where records are handed over as produced it appears instead as throughput sitting below the input rate.
- When can the widening not simply be moved after the movement?When the attached data is what the records are keyed on, or when a downstream filter reads its contents before the result narrows — then the records cannot be redistributed without it. The choice becomes carrying an identifier and re-associating on the collecting side, paying a lookup instead of the bytes, or accepting the movement and sizing the job for terabytes rather than for a record count.
- Does it matter whether the runtime would have pruned the unread fields itself?It matters enough not to rely on it. A declarative surface may push a projection below the exchange; a program written as opaque per-record functions gives the runtime nothing to see through, and the oldest model in this class rewrites nothing. Writing the projection yourself is correct on all of them, and costs nothing where the runtime would have done it anyway.
saying these in an interview costs you the question
- Estimates the movement from the record count and never checks record width.
- Assumes the runtime always drops unread fields before the move for you.
- Adds workers to a job whose time is bytes on the wire.
- Treats a step that multiplies record size as cheap, being one function call.
- Concludes the grouping is slow when the bytes reaching it are the cost.