skip to content

Why is the cost of redistributing data across a cluster counted in bytes moved rather than in records?

level: juniorimportance: must knowfreq 70%

answer

  1. count what the hardware moves
  2. three terms, not one
  3. wire, local disk, encode time
  4. width times record count

basics

~20 s

A redistribution costs bytes: what crosses the network, what the producing side writes to local disk where output is materialised, and processor time spent encoding it. Record counts predict none of those unless every record is the same width.

solid answer

~50 s

Cost a movement the way the hardware experiences it. A redistribution — every worker sending each record it holds to whichever worker will handle that record's key — spends three separate resources: bytes on the wire, bytes written to and read back from local disk on runtimes that materialise the producing side's output, and processor time spent encoding records into a byte stream and decoding them again at the other end. A record count is a proxy for none of those. Ten million narrow records can be far cheaper to move than a hundred thousand records each carrying a large document. So the estimate is average encoded width times record count, and I check the width first, because it is the number nobody measures. Two runs with identical record counts can differ by two orders of magnitude in bytes.

go deeper

for a junior

Recall that a movement is costed in bytes, not records, and that the estimate is average encoded width times record count. Be ready to say why two jobs with identical record counts can differ enormously.

for a middle

Explain all three terms — bytes on the wire, bytes written and read back on local disk where output is materialised, and processor time spent encoding and decoding — and say which of them a given runtime actually charges.

for a senior

Show that you measure width before you change anything: name the number you would establish first, and say what you would conclude differently if the same byte total were spread over far more, far smaller transfers.

for a principal

Frame it as which resource the platform is short of. Bandwidth, local disk and processor time are bought separately, and a fleet chosen for cores will bottleneck on movements long before a fleet chosen for network does.

## The unit of account A **redistribution** — the movement many engineers call a *shuffle*: every worker sends each record it holds to whichever worker will handle that record's key, so a step that needs records held on other machines can run at all — is usually the most expensive thing a distributed job does. Costing it means naming what it actually consumes, and it consumes three resources, not one: - **Bytes on the wire.** Every record that has to change machines crosses the network once, in whatever encoded form the runtime puts it in. This is the term most people mean by "what the movement costs", and it is bought from the cluster's bandwidth. - **Bytes on local disk.** Where the producing side writes its records into one bucket per destination and the collecting side comes back for them afterwards, each byte is also written to a local disk and read back off it. That is two disk operations per byte *on top of* the transfer. - **Processor time spent encoding and decoding.** Records that live as objects inside one worker's memory must be turned into a byte stream to travel and turned back into usable values at the far end. On a job whose real computation is a sum or a filter, this conversion can be most of the processor work the job performs. ## Why a record count predicts none of them A record count only stands in for bytes when every record is the same width, and outside a fixed-width table that is rare. Compare two movements: - ten million records of about 200 bytes each move roughly **2 GB**; - one hundred thousand records each carrying a 2 MB document move roughly **200 GB**. The second has one per cent of the record count and a hundred times the bytes. Any reasoning that starts "it is only a hundred thousand rows" has already lost. Width is set by things a record count cannot see: free text, nested structures, repeated groups, an attached payload, a wide set of fields where the step reads three of them. Three different widths are routinely confused, and only one of them is what a movement charges: | the number | what it is | where it matters | |---|---|---| | compressed size in the shared store | how large the input is in the durable storage every worker can reach | reading the input, not moving it | | in-memory footprint after decoding | how much of a worker's memory the same records occupy as live values | a worker's memory budget, which is a different subject | | **encoded bytes crossing the exchange** | what is actually serialised and sent, per record, times the records that move | **this is the movement's cost** | The gap between the first and the second is often several-fold; the gap between the first and the third depends on whether the runtime compresses what it sends. Quoting one of them when you mean another is the classic way an estimate ends up wrong by an order of magnitude. ## What varies between engines, and what does not All three terms exist in principle; which ones a given runtime actually charges you differs, and a good answer says so rather than describing one engine as the class: - Where a job is cut at **stage barriers** — a line across the job at which no downstream worker may compute until every upstream piece of the producing step has finished — the producing side typically writes its per-destination output to local disk, so the local-disk term is real and is paid twice per byte, once writing and once reading. - Where the runtime instead hands each record downstream as soon as it is produced, nothing is materialised: the local-disk term largely disappears, the wire and encoding terms remain, and per-message framing overhead becomes more visible because the units are small. - Some runtimes hold records in a compact binary layout they manage themselves; others hold ordinary language objects. The encoding term, and the memory footprint behind it, differ substantially between the two, which is why a number measured on one engine does not transfer. - Some runtimes compress intermediate output by default and some do not, so "bytes on the wire" may or may not already be the compressed figure. ## Estimating one before you run it 1. Take the number of records that reach the step needing the movement. 2. Take the average **encoded** width of only the fields that actually cross — not the in-memory footprint, not the source file's compressed size. 3. Multiply. Compare the product against what the cluster's network can carry in the time you are willing to wait. 4. Where the producing side materialises its output, add the same volume again as a local write and a local read, and expect the encoding work to scale with the same number. ## What this number is not - It is not money. Turning resource use into a bill and attributing it to a team is a separate subject. - It is not a planner's cost model. The number a query service computes internally to choose between strategies belongs to the warehouse's own execution, not here. - It is not rows, and the single most useful habit on this subject is refusing to let anyone state a movement's size in them.

  • Where does the local-disk term go on a runtime that hands each record downstream as it is produced?
    It largely disappears. Where a job is cut at barriers, the producing side writes its per-destination output to local disk and the collecting side reads it back, so each byte is written, read and sent. Where operators pass records on as they exist, nothing is materialised: the wire and encoding terms remain, but there is no local write — which is also why losing a worker means something different there.
  • Two movements carry the same total bytes, yet one takes much longer. What else varies?
    How the bytes are shaped. The same volume split into a few large transfers behaves very differently from the same volume split into millions of tiny ones, where per-transfer and per-record framing overhead dominate. Encoded width per record also changes how much processor time is spent for the same byte total, and a compressible payload and an incompressible one move different numbers of real bytes for the same logical size.

Freight is quoted by weight and volume, not by how many objects are in the crate. A pallet of feathers and a pallet of bricks have the same item count and completely different bills, and the invoice never asks how many feathers there were.

saying these in an interview costs you the question

  • Says a movement's cost is the number of records it moves.
  • Assumes every record is the same size, so counts stand in for bytes.
  • Ignores the processor time spent encoding and decoding records.
  • Believes every engine writes the producing side's output to local disk first.
  • Counts the network only, forgetting the bytes read back on the collecting side.