skip to content

When does compressing a step's output before it crosses the network make a redistribution slower rather than faster?

level: middleimportance: should knowfreq 48%

answer

  1. a trade, not a free win
  2. cores against bandwidth
  3. ratio near one buys nothing
  4. both sides pay to encode

basics

~20 s

Compressing trades processor time for bytes on the wire, so it loses whenever the job is short of cores rather than bandwidth, or when the data is already compact, so the ratio approaches one and the encoding work buys almost nothing.

solid answer

~50 s

Compression is a trade, not a free win: you spend processor time on the producing side to put fewer bytes on the network, and spend more processor time on the collecting side turning them back. It pays when the link is the constraint and the data actually compresses. It loses in two situations. First, when the cluster is already saturating its cores — then the encoding work lands directly on the critical path and the bytes you saved were never what was slowing you. Second, when the payload is close to incompressible — random identifiers, hashes, encrypted fields, an already-compressed document — so the ratio sits near one and both sides pay for nothing. Where the producing side's output is materialised to local disk, the saving lands twice, on the local writes and reads as well as the wire, which shifts the balance towards compressing.

go deeper

for a junior

Recall the trade in one line: compression spends processor time to move fewer bytes. It only helps if bytes were the problem in the first place.

for a middle

Explain both sides of the ledger — the ratio actually achieved on this data against the encoding and decoding time charged on the producing and collecting sides — and name a case where the ratio sits near one.

for a senior

Judge it from the job in front of you: say which resource is saturated before recommending any change, and note that where output is materialised the saving lands on local writes and reads as well as on the wire.

for a principal

Treat the default as a fleet-shaped decision rather than a per-job flag. What the machines and the links cost, and whether the same intermediate bytes are read more than once, decide it for a platform rather than for one pipeline.

## The trade, stated plainly A redistribution — every worker sending each record it holds to whichever worker will handle that record's key — spends bytes on the network and processor time turning records into those bytes and back. Compressing the output before it is sent moves cost from the first budget to the second. That is the whole of it: **fewer bytes on the wire, more processor time on both sides**. Whether it is a win therefore depends on which of the two the job is actually short of, and no general answer exists. ## When it pays - The network link, not the cores, is what the movement is waiting on. Workers with idle processor capacity and a saturated link are the textbook case. - The data genuinely compresses. Repetitive text, low-cardinality string fields repeated across millions of records, sparse or defaulted fields — all shrink substantially. - The producing side's output is **materialised**: written into one bucket per destination on local disk and collected afterwards. Then compression saves the local write and the local read as well as the transfer, so the same processor spend buys three savings instead of one. - The same bytes are read more than once. Paying the encoding cost once and the decoding cost several times can still come out ahead. ## When it loses 1. **The job is already processor-bound.** If every core is busy, compression's work does not overlap with anything; it extends the producing side directly, and the bytes it removed were not the constraint. 2. **The ratio is near one.** Random identifiers, hashes, encrypted values, and payloads that arrived already compressed do not shrink. Both sides still pay to run the codec over them. Re-compressing an already-compressed document is not impossible, as is sometimes claimed — it simply returns almost nothing while charging full price, and can even grow the stream slightly. 3. **The records are tiny and numerous.** Per-block framing and per-message overhead can exceed the saving when each unit is a few dozen bytes. 4. **The choice was made for maximum ratio on a fast network.** Ratio and speed are separate axes, and past some point extra ratio costs more processor time than the removed bytes were worth. ## Where the cost lands | side of the movement | what it pays | when it shows up | |---|---|---| | producing side | encoding each record, then compressing the bucket it writes | while the producing step is still running, directly on its wall-clock | | the network | fewer bytes, in proportion to the ratio actually achieved on this data | during the transfer | | collecting side | decompressing, then decoding every record it reads | as it pulls its share, before any of the downstream work begins | The important point in the table is that the saving is one row and the cost is two. A candidate who only names the middle row is reasoning about half the ledger. ## What varies between engines Describe the mechanism, not one product's behaviour: - Some runtimes compress intermediate output by default and others do not, so "the bytes on the wire" may already be the compressed figure — or may not be — before you change anything. - Where the runtime cuts the job at barriers and materialises output, the local-disk saving is real; where it hands each record downstream as it is produced, there is no local write to save and only the wire term is in play. - Runtimes that hold records in a compact binary layout they manage have already removed much of the redundancy that a general-purpose codec would otherwise find; runtimes holding ordinary language objects produce a byte stream with more slack in it, and compress better for the same effort. ## How to decide in an interview answer Say which resource is scarce, and say how you would know. The honest sequence is: establish the byte volume of the movement and the processor utilisation of the workers while it runs; if cores are idle and bytes are large, compression is the lever; if cores are pinned, look at removing bytes instead of re-encoding them — dropping fields nobody downstream reads, or not widening the record before the move at all. Which particular codec trades ratio against speed how is a serialization question rather than a movement question, and a good answer hands it off rather than reciting it. The failure mode to avoid is treating a compression setting as a performance knob with an obviously correct position. It has two positions, each correct for a different cluster, and the data in front of you decides.

  • If the data is compressible and the cores are free, is more compression always better?
    No. Ratio and speed are separate axes, and past some point extra ratio costs more processor time than the bytes it removes are worth. Choose by which resource is scarce and by how large the movement is; how particular codecs sit on that ratio-against-speed curve is a serialization subject rather than a movement one.
  • Where does the processor cost actually land, and on which side of the movement?
    Encoding and compressing are charged to the producing side as it writes its output; decompressing and decoding are charged to the collecting side as it reads that output back. So one decision shows up as processor time in two different places in the job, and a cluster already saturating its cores pays both in wall-clock rather than absorbing either.

saying these in an interview costs you the question

  • Says compression always speeds a movement because fewer bytes travel.
  • Chooses the highest ratio available without asking what the cores are doing.
  • Assumes every payload compresses by roughly half regardless of its content.
  • Forgets the collecting side pays to decompress and decode everything sent.
  • Treats it as free because the runtime does it somewhere in the background.