skip to content

Compute Near the Data

Placing a piece of work on the machine that already holds its bytes, the rungs it falls back through when that machine is busy, and why separating storage from compute traded most of it away.

on this pageshow

questions

4

In a cluster where storage and compute share machines, why place a piece of work on the machine already holding its bytes?

level: juniorimportance: must knowfreq 58%

answer

  1. small thing moves, large thing stays
  2. storage and compute on one machine
  3. local disk read against a shared switch
  4. gain equals the local-remote read gap
  5. input reads only, not redistribution

basics

~20 s

The program is kilobytes and the input is gigabytes, so shipping the work to the bytes is far cheaper than pulling the bytes across a shared network to wherever a thread happens to be free.

solid answer

~50 s

The idea belongs to one arrangement: storage running on the same machines as the compute, with each block of a file on the local disks of a few of them. A piece of the input is one slice that a single worker thread reads start to finish, so before it computes anything the bytes must reach its memory — either from a disk attached to its own machine, or over a network link shared with every other machine. Since the program is tiny and the input is enormous, the coordinating process that turns the program into a graph and hands out the pieces prefers to send the work to the bytes. The benefit is exactly the gap between a local read and a remote one, which is why the idea faded as networks got faster and storage moved off the compute machines.

go deeper

for a junior

Recall the ratio: the program is tiny, the input is huge, so the work travels to the bytes. Be able to say that this only makes sense when the machines that compute are also the machines that store.

for a middle

Explain what the saving actually is — the difference between a disk read on the same machine and a read over a shared network link — and that it scales with bytes read, not with cluster size.

for a senior

Show that you know the scope: placement affects input reads only, never the redistribution a regrouping step forces, and the whole optimisation is absent when storage is remote from compute.

for a principal

Frame it as an architecture-era trade. The preference was worth a scheduler's complexity when networks lagged disks; treat any platform still paying for that complexity over remote storage as carrying dead weight.

## What "near the data" actually means A **piece of the input** is one slice of the stored input that a single **worker thread** — one lane inside a worker process, which is itself one operating-system process on one machine — reads and processes from start to finish. Before that thread can compute anything, the slice's bytes have to be in its memory, and there are only two routes: - **Local**: the bytes already sit on a disk attached to the machine the thread runs on, and the read never leaves the box. - **Remote**: the bytes sit on some other machine, and they cross a network link that every other machine in the cluster is also using. Placing compute near the data is the scheduling preference that tries to produce the first case. The **coordinating process** — the single process that turns the program into a graph, decides the pieces and hands them out — asks where each slice's bytes live and tries to hand that slice to a worker thread on one of those machines. ## The arrangement that made it pay The idea was invented for a **cluster file system**: storage running on the same machines as the compute, where each block of a file is written to the local disks of a small number of those machines. That arrangement is what gives a scheduler a choice worth making — there exists a machine that *has* the bytes, so preferring it means something. On storage that no compute machine holds, the preference has nothing to express at all. Against that background the arithmetic is one-sided: | what moves | typical size | how often | |---|---|---| | the program | kilobytes to a few megabytes | once, to each machine that runs part of it | | the input | gigabytes to terabytes | every byte, every run | Shipping a small program to where the large data already sits costs almost nothing. Shipping the data to wherever a lane is free costs the whole input, over a path shared with every other piece of work in the cluster. The old slogan — *move the computation, not the data* — is a claim about that ratio, not a law. ## What the saving is, precisely Interviewers push here, and the precise answer is worth having: 1. **It is the gap between a local read and a remote read, and nothing more.** If local disk and the network deliver comparable bandwidth for this input, a local placement saves nothing you can measure. 2. **It is per-byte and per-run.** A run reading ten terabytes benefits ten times as much as one reading one terabyte; the program shipped in the other direction is the same size either way. 3. **It is congestion relief as much as latency.** Hundreds of threads reading remotely converge on the same switches. Local reads keep that traffic off the shared path entirely, which helps every other job on the cluster too. 4. **It applies to reading the input, and only to that.** Once the program reaches a step a worker cannot finish from the records it already holds, records have to be redistributed between workers whatever the placement was. These are two different costs: reading remotely happens *before* the work, redistribution happens *because of* it, and placement only influences the first. ## What it is not - It is **not** a copy step. Nothing relocates bytes to a chosen machine before the run; the scheduler chooses among machines that already hold them, or gives up and reads across the network. - It is **not** a correctness mechanism. A remote read returns the same bytes; only the cost differs. - It is **not** the same word as the on-disk grouping of files into directories named for a column's value. That grouping decides which files are read at all; placement decides which machine reads them. ## Where engines genuinely differ Do not state one behaviour for the whole class: - Some runtimes model an explicit ordered preference and will briefly hold a piece hoping for a nearer machine; others express no preference and place purely by which lane is free. - In a **finite job** — a run over an input that ends, so the bytes can be measured before it starts — the pieces are derived from stored data and can be matched to the machines holding it. In a **continuous job** — a run over an input with no end — there is nothing stored to match against; operator placement is settled when the job starts rather than re-decided per record. - Where the input lives in an **object store** — remote blob storage reached over the network — no machine holds the bytes, every candidate is equidistant, and the preference is simply absent. That is the common modern deployment. ## The answer that lands Name the three conditions the optimisation needs: compute sharing machines with storage, an input far larger than the program, and a network materially slower than local disk. Someone who can name those has also explained, without being asked, why most deployments today get nothing from it.

  • Does a local placement help a step that has to regroup records by key?
    No. A step that needs records currently sitting on every other worker forces those records to travel regardless of where the pieces were placed. Placement only affects how the input bytes reached a machine in the first place; the regrouping cost is a separate movement with its own accounting, and no placement choice removes it.
  • How would you tell whether a run actually got local reads?
    Runs that model a preference usually report, per piece, whether the bytes were read from a disk on the same machine or fetched across the network; compare that count against the total piece count for the run. Where the input is remote storage, the counter is meaningless because no placement was ever nearer than another.
  • If the network is as fast as local disk, is the preference harmful?
    Not harmful, but not free either. Expressing a preference means a scheduler may leave a lane idle while it looks for a nearer machine, so on a cluster where local and remote reads cost the same, that idle time is pure loss. The correct posture there is to place by free capacity.

A researcher and a reference library. Flying the researcher to the library costs one seat; shipping the library to the researcher costs a fleet of trucks. The moment the library is digitised and reachable from anywhere at the same speed, the seat stops being worth booking.

saying these in an interview costs you the question

  • Thinks the bytes are copied to a chosen machine before the work starts
  • Says a remote read returns different or stale data than a local read
  • Believes the preference still applies when the input is in remote storage
  • Claims local placement removes the need to redistribute records later
  • Assumes network bandwidth inside one datacentre is effectively unlimited
  • Treats the optimisation as a correctness rule rather than a cost preference
open as a page

When the machine holding a piece's bytes is busy, should the scheduler wait for it or start elsewhere?

level: middleimportance: should knowfreq 44%

basics

~20 s

It depends on how long the piece will run. Waiting buys a cheap local read but burns an idle lane now; starting further away pays a remote read for the piece's whole life. Short pieces should never wait long.

open as a page

When a job's input lives in remote blob storage, what happens to the scheduler's preference for nearby bytes?

level: seniorimportance: should knowfreq 52%

basics

~20 s

It disappears, on purpose. No compute machine holds the bytes, so every candidate is equally far and placement collapses to free capacity. The industry traded that optimisation for compute it can resize or discard independently of the data.

open as a page

When every input read is remote, what work inside a running job is still tied to one particular machine?

level: seniorimportance: nice to knowfreq 32%

basics

~20 s

Two things: bytes a machine wrote to its own disk during the run, and work for a key whose accumulated state one worker already holds. Both make that machine non-interchangeable, which is the real cost.

open as a page