skip to content

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%

answer

  1. input equidistant, run-made data is not
  2. producer's own disk, and the key's owner
  3. the run manufactures its own locality
  4. affinity fights disposability
  5. record-at-a-time engines have less of it

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.

solid answer

~40 s

Input placement dies when the input is remote, but nearness reappears inside the run in two places. First, bytes a machine produced for itself — intermediate results held on its own local disk — are cheap for it to read back and expensive for anyone else, so the work consuming them wants to be there. Second, work for a key whose accumulated state one worker already holds must reach that worker, or the state has to be fetched or rebuilt first. Both are genuine nearness and both have the same consequence: that machine is no longer interchangeable with any other. Losing it means recomputing or restoring what it held, and removing it deliberately means moving that first. Affinity and elasticity pull against each other, and how much of each exists differs sharply between runtimes.

go deeper

for a junior

Recall that a job creates data of its own while it runs, and that data sits on a particular machine even when the input it read came from far away.

for a middle

Explain the two cases: bytes a machine wrote for itself, and accumulated values for a key that one worker owns. Both mean the work has a nearer place to run.

for a senior

Draw the consequence. Affinity makes a machine non-interchangeable, so losing or reclaiming it sets the run back, and say which of the two survivors the runtime you are describing actually has.

for a principal

Weigh it as a platform choice. Intermediate data can be pushed off the compute machines too, buying disposability back at the price of more network traffic; decide which of your workloads is worth that.

## Two survivors Once the input is remote blob storage, no machine is nearer to the input than any other and the scheduler's input preference is vacuous. But a running job produces data of its own, and that data does have a home. **1. Bytes the machine produced for itself.** Many runtimes hold intermediate results — output a later step will consume — on the local disk of the worker process that produced them. Reading those back on the same machine is a local read; reading them from anywhere else means a fetch over the network. So the work that consumes them has a genuinely nearer candidate, and that candidate is one specific machine. **2. Work bound to a key.** Where a job accumulates something per key, the accumulated value has to live with one worker, the one that owns that key. A record for that key arriving anywhere else either has to be routed to the owner or has to have the state fetched or rebuilt before it can be processed. That routing is a nearness constraint in exactly the same sense as the original one, and the mechanics of key-bound state are owned by the long-lived-state subject. Notice what both have in common and how they differ from input placement: | | input bytes (remote store) | intermediate bytes | key-bound state | |---|---|---|---| | is any machine nearer? | no | yes, the producer | yes, the key's owner | | created by | storage, before the run | the run itself | the run itself | | survives machine loss? | yes, untouched | no, must be recomputed or refetched | only via a saved recovery point | ## Affinity is the opposite of interchangeability This is the part worth saying out loud in an interview. The reason remote storage was worth its bandwidth cost is that compute machines became interchangeable and therefore disposable. Every scrap of nearness a run creates for itself claws some of that back: - A machine holding intermediate bytes cannot be reclaimed for free. Whatever depended on those bytes has to be recomputed from its inputs, or the bytes fetched elsewhere first. - A machine holding key-bound state cannot be removed without moving that state, and it cannot be replaced by a fresh machine without restoring it. - Adding machines mid-run does not immediately help the work that is pinned, because the pinned work cannot follow. How far this goes — what is recomputed, what is restored, how machines are added or dropped mid-run — belongs to the failure-handling and machine-supply subjects. What belongs here is the cause: the run manufactured its own locality, and locality is what makes a machine special. ## Where runtimes genuinely differ This is the claim most likely to be wrong if you generalise from the engine you know: - **Batch-lineage runtimes** typically materialise intermediate results to the producing machine's own disk, so the first survivor is strongly present. **Record-at-a-time continuous runtimes** often push records across the network as they are produced with no materialisation at all, so they have little or none of it — their only affinity is state. - The **oldest model in this class**, which writes every intermediate result to disk between the two halves of a job, has the first survivor in its most extreme form, and it is exactly why that model's recovery story is simple. - Where continuous work is run as a **rapid succession of small finite jobs**, intermediate data is short-lived and the per-key accumulation is kept somewhere the next small job can pick it up, which makes the affinity look different again. - Some deployments store intermediate results away from the compute machines on purpose — the same trade as the input, made again — which removes the first survivor deliberately so that machines stay disposable. ## How this shows up in practice 1. **A machine is reclaimed late in a long run and the run takes a large step backwards.** The bytes it held were not input; they were produced, and nothing else has them. 2. **A job is restarted at a different width and the first records are slow.** Key ownership moved, so accumulated state had to be re-established before the work could proceed. 3. **A cluster scaled out mid-run and throughput did not rise proportionally.** The pinned work could not move to the new machines. ## The interview form Asked "does data locality still matter?", the answer that separates a senior candidate is: not for reading the input, which is now equidistant from everywhere by design, but yes for what the run makes — and the price of that surviving nearness is that some machines stopped being disposable. Then say which of the two survivors your runtime actually has, because they are not the same in every engine.

  • What does this affinity cost when a machine is taken away mid-run?
    Whatever that machine alone held has to be reproduced. Intermediate bytes must be recomputed from their inputs or fetched from a copy if one exists, and key-bound state must be restored from a saved recovery point and re-owned by another worker. How that reproduction is scoped is the failure-handling subject's business; the reason it is needed is the affinity.
  • Does a continuous job have both survivors?
    Not necessarily. A record-at-a-time runtime that pushes records across the network as it produces them keeps little or no intermediate data on local disk, so its only real affinity is key-bound state. A runtime that instead runs continuous work as repeated small finite jobs has a short-lived version of both. Say which model you mean before making the claim.
  • Can a platform choose to give up the intermediate-bytes affinity as well?
    Yes, and some do: write intermediate results to remote storage or a separate service instead of the producing machine's disk. It is the same trade made a second time — slower reads and more network traffic, bought in exchange for machines that stay interchangeable and can be reclaimed without setting the run back.

saying these in an interview costs you the question

  • Says no nearness at all survives once input is in remote storage
  • Claims every engine keeps intermediate results on the producer's local disk
  • Thinks a pinned machine can be swapped for a fresh one at no cost
  • Assumes adding machines mid-run immediately speeds up pinned work
  • Treats state affinity and input placement as the same mechanism