skip to content

Why is copying a 2 GB input to every worker more expensive as the cluster grows, when the input never changes size?

level: middleimportance: must knowfreq 60%

answer

  1. per worker, not per job
  2. multiply by cluster width
  3. resident everywhere at the same moment
  4. more workers means more copies

basics

~20 s

The copy is paid once per worker, not once per job: two hundred workers means 2 GB sent two hundred times and 2 GB resident on two hundred machines at once. Cost tracks cluster width.

solid answer

~50 s

Copying an input to every worker - a worker being one process on one machine running some of the job's pieces - replicates it, so the input's own size is only one factor. The bytes on the wire are that size multiplied by the number of workers: 2 GB across 200 workers is 400 GB moved. The second cost is memory: the same 2 GB is resident on all 200 workers simultaneously for as long as the step needs it, which is 400 GB of cluster memory doing one job. Both grow with width, so widening a cluster to speed a slow job makes the copy worse, not better. How the copies are delivered varies - some engines gather the input and push it out, others have each worker read it from shared storage, others pass it worker to worker - and that changes the delivery time, not the arithmetic.

code

python · 21 lines
python
# Judging a copy-to-every-worker plan by arithmetic alone.
# Sizes in bytes. The DECODED size is what each worker must hold.

def copy_cost(decoded_bytes, workers):
    return {
        "on_the_wire": decoded_bytes * workers,
        "resident_per_worker": decoded_bytes,
        "resident_cluster_wide": decoded_bytes * workers,
    }

decoded = 2_000_000_000        # 400 MB stored, about 2 GB once decoded
larger_input = 900_000_000_000 # the input we are trying not to move

for workers in (50, 200, 800):
    c = copy_cost(decoded, workers)
    print(workers, c["on_the_wire"], c["resident_per_worker"])

# 50  -> 100 GB on the wire, 2 GB resident on each worker
# 200 -> 400 GB on the wire, 2 GB resident on each worker
# 800 -> 1.6 TB on the wire, which is now a sizeable fraction
#        of simply moving the 900 GB input once

go deeper

for a junior

Remember that the copy goes to every worker, so the number of machines is part of its price. The input being small does not make the copying small.

for a middle

Do the multiplication out loud - size times worker count for the wire, and the same size resident on every worker at once - and note that widening the cluster grows both.

for a senior

Show where the second cost hides: monitoring reports bytes transferred, not the same data resident in hundreds of places, so the memory side of this is usually found by failure rather than by a graph.

for a principal

Treat it as a platform question: a shape that gets worse as the cluster grows caps how far a job can be scaled out, and the resident memory is capacity denied to every other job on those machines.

## The arithmetic, first Every other statement about this technique follows from one line: **the copy is paid once per worker, not once per job.** A worker here is one process on one machine that holds a slice of the job's memory and runs some of its pieces. Take a smaller input that occupies 2 GB once it is decoded into records in memory, and a step that matches it against a much larger input left where it was read: | Workers | Bytes on the wire for the copy | Resident per worker | Cluster memory held at once | |---|---|---|---| | 50 | 100 GB | 2 GB | 100 GB | | 200 | 400 GB | 2 GB | 400 GB | | 800 | 1.6 TB | 2 GB | 1.6 TB | The input never changed. The cluster did. ## Two costs, and both are multiplied - **Bytes on the network.** The copy is delivered to every worker, so the total moved is the copy's size times the number of workers. It is the one movement cost in a job that grows when you add machines. - **Memory held simultaneously.** Every worker holds the whole copy at the same time, for as long as the step that needs it runs. Unlike a redistribution, where each worker receives only its own share, there is no division of labour here: replication is the point. The second cost is the one candidates forget, because monitoring usually shows bytes transferred and rarely shows "the same object, resident in four hundred places". ## Per worker, or per machine? This genuinely differs and is worth asking about rather than asserting: - Some engines hold one copy **per worker process**, so three workers on one machine hold three copies and the machine pays three times. - Others cache one copy **per machine** and let the workers there share it, which cuts the resident cost by the number of workers per machine but not the wire cost of getting it there once. The arithmetic above is the per-process case, which is the conservative one to plan against. ## How the copies get there The delivery mechanism also varies, and it changes how long the distribution takes without changing the total bytes: 1. **Gathered, then pushed.** The input is pulled back to the process coordinating the job, which sends it out. Simple, and the sending process becomes a bottleneck and a single point of memory pressure as the cluster widens. 2. **Read from shared storage.** Each worker fetches the input itself from durable storage every worker can reach. The fan-out is the storage system's problem, not one process's. 3. **Passed onward between workers.** Workers that already hold the copy serve it to those that do not, so distribution time grows roughly with the logarithm of the cluster width rather than linearly. All three move the same number of bytes in total. Only the time to finish distributing, and where the strain lands, differ. ## Why adding workers can make this job slower The intuitive move when a job is slow is to widen the cluster. For a step matching against a copied input, widening does two things at once: it divides the larger input into more pieces, which helps, and it multiplies the copy, which does not. Past some width the distribution of the copy and the memory it occupies dominate, and the job stops improving or gets worse - and it fails outright at the width where the copy no longer fits beside everything else the worker is doing. That is the crossover to hold in your head: this technique is cheapest on a **narrow** cluster with a genuinely **small** input, and it degrades along both axes at once. ## How long the copy is held Once a worker holds the copy, further steps in the same job can match against the same resident data without it being sent again. When it is released differs: some engines free it when the step that needed it completes, others keep it until the program releases it or the job ends. A copy nobody released is memory nobody else can use, on every worker, for the rest of the run - which is the same arithmetic again, now charged for a duration rather than a moment. ## What to say in the interview Give the multiplication out loud. "Two gigabytes decoded, four hundred workers, so eight hundred gigabytes on the wire and eight hundred gigabytes of cluster memory held at the same moment - the input is small, the copy is not."

  • Does the copy have to be sent again for every step that uses it?
    Usually not: once a worker holds it, later steps in the same job can match against the resident copy. When it is released varies - some engines free it as the step ends, others hold it until the program releases it or the job finishes - so a copy that is never released occupies memory on every worker for the rest of the run.
  • The cluster has four workers per machine. How does that change the memory arithmetic?
    It depends on whether copies are held per worker process or cached once per machine and shared. In the first case that machine holds four copies, in the second one. Plan against the per-process figure unless you know otherwise, because the difference is fourfold.

Handing every person in a meeting their own printed copy of the same one-page memo. The memo does not get longer when more people come, but the printing, the carrying and the desk space are paid per attendee - so a bigger meeting costs more even though nothing about the document changed.

saying these in an interview costs you the question

  • Says the copy costs its own size once, however wide the cluster is.
  • Thinks adding workers always shortens a job that copies an input everywhere.
  • Counts network bytes only and ignores memory resident on every worker.
  • Believes the copies are always sent sequentially by one coordinating process.
  • Assumes the copy is released the instant the first match completes.