skip to content

Splitting the Work

How the work is cut into pieces and how many, where each piece runs and what shape it leaves behind: the choices that set how much of the cluster works at once, and who gets to make them.

on this pageshow

explore

questions

28

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

A job's only input is one 40 GB file compressed as a single stream — how many worker threads can read it?

level: juniorimportance: must knowfreq 68%

basics

~20 s

One. A stored form that cannot be decoded from an arbitrary offset yields exactly one piece of the input — one slice a single worker thread reads end to end — so every other thread in the cluster stays idle during that read.

open as a page

A worker thread holds one piece of the input. What decides whether it can finish a step alone or needs records other workers hold?

level: juniorimportance: must knowfreq 78%

basics

~20 s

Dependency shape decides: if each output depends only on records already in the thread's piece, it finishes alone — a narrow step. If one output needs records spread across every piece, the step is wide.

open as a page

A finite job's stored input is divided into eight pieces and the cluster offers 200 worker threads; how many can work at once?

level: juniorimportance: must knowfreq 80%

basics

~10 s

Eight. The piece count is the ceiling on how many worker threads can be busy; the machine count only decides how many threads exist, so the other 192 lanes stay idle.

open as a page

A finite job's final step holds 2,000 pieces and writes to a directory. How many files does it leave, and what set that number?

level: juniorimportance: must knowfreq 70%

basics

~20 s

Roughly one file per piece, so about 2,000. Each worker thread writes the records it already holds, where it is; nothing gathers them first. The piece count at the final step, not the machine count, set the file count.

open as a page

Which two quantities set the piece count for a finite job reading 800 GB on 400 worker threads?

level: middleimportance: must knowfreq 74%

basics

~20 s

A target size per piece and the total worker-thread count. Size keeps each piece inside one thread's memory and worth its fixed cost; thread count keeps the lanes busy in even waves. Reconcile them rather than using either alone.

open as a page

Who fixes the piece count of a distributed job, and when - at submission, mid-run, or only on restart?

level: middleimportance: must knowfreq 66%

basics

~20 s

It varies by engine and by job kind. Some derive and freeze the count at submission; some let the author set it per redistribution point and adjust it between step groups while running; a continuous job's declared width usually stands until a restart, and some platforms never expose it at all.

open as a page

Stored files are grouped into directories named for one column's value — what does a filter on that column change before reading?

level: middleimportance: must knowfreq 72%

basics

~20 s

Whole directories are matched against the filter while the plan is being built, so their files never become pieces of the input at all. Both the piece count and the bytes read fall, before a single file is opened.

open as a page

In a pipeline that filters rows, derives a column, counts distinct users per country and orders the result, which steps regroup records?

level: middleimportance: must knowfreq 70%

basics

~20 s

The filter and the derived column are narrow: each output depends only on the record in hand. The per-country distinct count is wide, and so is the global ordering — their outputs depend on records that may sit in any piece.

open as a page

How does a finite job arrive at its piece count, and a continuous job at its declared operator width?

level: middleimportance: must knowfreq 66%

basics

~20 s

A finite job derives it: the stored input is measured and cut into pieces before the run. A continuous input has no end to measure, so the author declares how many copies of each operator run.

open as a page

A write grouped into directories by a 200-value country column runs over 2,000 pieces. Why can it leave 400,000 files?

level: middleimportance: must knowfreq 62%

basics

~20 s

Because the two numbers multiply. Each worker thread writes what it holds into the directory each record's value names, so a piece holding all 200 values leaves 200 files. Across 2,000 pieces that is 2,000 x 200 files, each tiny.

open as a page

Why does cutting a 2 GB input into 50,000 pieces usually run slower than cutting it into 200?

level: juniorimportance: should knowfreq 60%

basics

~20 s

Each piece carries a fixed cost that does not shrink with its bytes: it is created, dispatched, tracked and collected, and its input is opened and closed. At 40 KB per piece that bookkeeping outweighs the reading, and the coordinating process becomes the bottleneck.

open as a page

Two nightly datasets are joined on customer id; what must be true of both writes for that join to move no records?

level: middleimportance: should knowfreq 45%

basics

~20 s

Both writes must have used the same key, the same function from key to piece number, and piece counts the engine can relate - normally the same count. Only then does every matching pair already sit on one worker.

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

A job groups records by account, then joins on account, then aggregates by account again. How can it be rewritten to regroup once?

level: middleimportance: should knowfreq 56%

basics

~20 s

Regroup once on the account key, then keep that division. After the first wide step every account's records already sit together, so later steps keyed on the same account finish locally — unless something in between re-divides or hides the key.

open as a page

A job's 30 GB input is stored in a form that yields exactly one piece — what does adding machines buy?

level: middleimportance: should knowfreq 58%

basics

~20 s

No extra concurrency. One piece is read by one worker thread from start to finish, so the run's ceiling is one lane and the rest of the cluster idles no matter how many machines are added.

open as a page

A nightly job's piece count was sized when its input was 200 GB and it now reads 3 TB unchanged - what rots, and how would you tell?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Bytes per piece grew fifteenfold while the count stayed put, so every piece is now oversized: uniformly longer runtimes, new disk spilling, memory pressure and expensive retries. The tell is that the whole duration distribution shifted, not that one piece is slow.

open as a page

A join between two tables stopped avoiding redistribution six months after both were written to match - what eroded the match?

level: seniorimportance: should knowfreq 33%

basics

~20 s

Later writes that did not follow the agreement: a backfill or a second pipeline appending files divided some other way, a piece count re-tuned on one side, or a key column quietly redefined. Nothing fails; the saving just disappears.

open as a page

Both join inputs were written with an identical division on the join key, yet the run still redistributes both sides - why?

level: seniorimportance: should knowfreq 42%

basics

~20 s

Because the run must know the layout holds, not merely benefit from it. Unless the key, rule and piece count are recorded where the plan reads them and survive every earlier step, it moves the records.

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

A job reads 300,000 files of 50 KB from an object store and idles twenty minutes before its first record — why?

level: seniorimportance: should knowfreq 62%

basics

~20 s

Planning dominates, not reading. Listing 300,000 objects and turning each into a piece of the input costs far more than the 15 GB of payload, and that work is serial, so idle worker threads cannot absorb any of it.

open as a page

Your program contains three wide steps. How does that number read differently on a runtime that schedules in units versus one where all operators run at once?

level: seniorimportance: should knowfreq 42%

basics

~20 s

On a runtime that schedules the work between regroupings, three wide steps mean four separately scheduled units of steps. On a runtime where every operator runs at once, they mean three routing hops each record may make in flight.

open as a page

Your input tripled: the nightly finite job now uses more of the cluster and the continuous job does not — why?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Because the two numbers have different origins. The finite run derives its pieces from the stored bytes, so more data means more pieces and more busy lanes; the continuous job runs at a width the author declared, which no change in the input can move.

open as a page

A finite job's 800 worker threads write rows straight into one live operational database at the final step, and the database stalls. What has to change?

level: seniorimportance: should knowfreq 54%

basics

~20 s

The limit is now the destination's, not the cluster's. Three things change: rows go in batches rather than one request each, the number of simultaneous writers is capped independently of the piece count, and the run holds a rate ceiling the store was measured to survive.

open as a page

A shared dataset's files sit in directories named for a high-cardinality column — how do you judge that granularity for every job that reads it?

level: principalimportance: should knowfreq 45%

basics

~20 s

Judge it against the population of reading jobs, not one query: how often each candidate column is filtered on, how many directories and how large a file each granularity leaves, and how skewed the column is. Fine granularity buys elimination for one access pattern and charges everyone else.

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

Forty jobs join one shared table on the same key - should the platform require every write to preserve a matching division, and what does that cost?

level: principalimportance: nice to knowfreq 26%

basics

~20 s

Only if that key dominates the joins and the table is read far more often than written. The standard freezes a piece count into shared data, taxes every writer, and needs an owner and a check to survive.

open as a page

Dozens of jobs across a platform write into the same shared operational stores. What house rules stop any single run from flattening one?

level: principalimportance: nice to knowfreq 36%

basics

~20 s

Make the write budget a property of the destination, not of each job: a measured concurrency and rate allowance per store, divided among the jobs that write to it, with batching required, writer concurrency decoupled from the piece count, and a defined behaviour when the allowance runs out.

open as a page