skip to content

The Coordinating Process

Every distributed job has one process that is not a worker: it plans the pieces, hands them out, tracks what finished, and receives anything the program asks to bring back to one place.

on this pageshow

questions

4

In a distributed job, what does the single process that is not a worker do, and what does it never do?

level: juniorimportance: must knowfreq 72%

answer

  1. one process is not a worker
  2. plans, hands out, tracks completion
  3. no piece of data runs in it
  4. returned records land in its memory
  5. sized by piece count, not input

basics

~20 s

One process, the coordinating process, turns the program into pieces of work, hands them to worker processes, tracks which finished, and receives anything the program asks to bring back. It runs no piece of the data itself.

solid answer

~50 s

A distributed job runs as many processes across several machines. Almost all are **worker processes**: the processes on cluster machines that actually run pieces of the job and own the memory those pieces use. Exactly one is not. The **coordinating process** plans the pieces, hands them out, tracks what finished and re-issues what failed, and receives anything the program explicitly asks to bring back to one place. It never reads the bulk of the input and never computes a share of the data — so size it from the bookkeeping and from what is returned, never from the input volume. How busy it is while the job runs varies by runtime: where continuous work is executed as a rapid succession of small finite jobs it plans a fresh wave every round, while a record-at-a-time runtime plans once and then mostly tracks health and recovery points.

go deeper

for a junior

Recall the division of labour: worker processes run pieces of the data, and one coordinating process plans the pieces, hands them out and tracks what finished. Say plainly that it computes no share of the data.

for a middle

Explain what actually lives in that process — the plan, the assignment table, the completion record and anything the program asked to bring back — and derive the sizing rule from it: piece count and returned size, never input volume.

for a senior

Show that you know how busy it is depends on the runtime, and connect that to operations: a per-round planner on the critical path of every round behaves very differently under a large piece count from one that plans once at start-up.

for a principal

The angle is where the line sits between the platform's supply and the team's job. On a managed compute service the coordinating process is the provider's to own and size, which removes a tuning lever and a failure the team would otherwise have to design around.

## Two kinds of process in one job A distributed job is not one program on one machine; it is a set of operating-system processes spread over several machines, all working on one submitted program. Almost all of them are **worker processes** — the processes on cluster machines that actually run pieces of the job and own the memory those pieces use. Exactly one of them is not a worker. That one is the **coordinating process**: the single process in a distributed job that plans the pieces, hands them out, tracks what finished, and receives anything the program asks to bring back to one place. This is a first-screen question because almost everything else about running work on a cluster depends on the split. An engineer who pictures the cluster as one uniform bag of machines cannot explain why one out-of-memory error ends the run while another costs thirty seconds, and cannot explain why adding machines sometimes changes nothing at all. ## What the coordinating process holds - **The plan.** The submitted program is turned into a graph of steps, and each step into a set of pieces that can be run independently. That structure lives here. - **The assignment table.** Which piece went to which worker process, when, and on which attempt. - **Completion and failure bookkeeping.** Which pieces reported success, which failed, which must be handed out again. On a job cut into tens of thousands of pieces this is itself measurable memory and measurable processor time. - **Whatever the program asked to bring back.** An operation that returns records to the submitted program moves them out of the workers and into this one process's memory. That is the usual reason it runs out of room. - **The submitted program's ordinary code.** Variables, control flow, the loop that decides what to compute next. Anything not expressed as work over the distributed data runs here, single-threaded, on one machine. ## What it does not hold It does not hold the input. A job whose input is far larger than any one machine still has a coordinating process that never sees most of it: the input is divided by describing which range of which file each piece covers, and the bytes are read by whichever worker process is assigned that piece. It does not hold the per-key working set an aggregation builds. It does not write the output — each worker writes its own piece to the destination. The practical rule that follows: **its memory tracks the number of pieces and the size of what is returned, not the size of the input.** A job over a hundred terabytes that returns one number needs no more coordinating memory than a job over a hundred megabytes that returns one number. ## Where engines differ How busy this process is while the job runs is a property of the runtime, not of the class: | how the runtime executes work | what the coordinating process does while running | |---|---| | continuous work executed as a rapid succession of small finite jobs | plans and hands out a fresh wave of pieces for each round of work, so it sits on the critical path of every round | | record-at-a-time, with long-lived key-bound state | plans once at start-up; records then flow directly between workers while it tracks health and coordinates recovery points | | a two-phase disk-to-disk model | plans the two phases and tracks piece completion, while a separate component owns the pool of machines itself | Two statements hold across all of them. First, the data records do not travel through the coordinating process on their way between workers; records moving between workers is an entirely separate path, and the only records that reach the coordinating process are the ones the program asked to have returned. Second, on a **managed compute service** — a service you hand work to and are never shown a machine by — this process still exists with all the same limits, but the provider owns it and you can neither see it nor size it. ## What an interviewer is listening for 1. **Planning and bookkeeping, not computation.** A candidate who says it also runs a share of the data has the model wrong. 2. **A different failure mode.** Losing a worker process costs the pieces it was running; losing the coordinating process leaves nothing planning, handing out or tracking. 3. **A different sizing input.** Piece count and returned size, never input volume. 4. **One process on one machine.** The parts of the program outside the distributed operations do not get faster when the cluster grows.

  • A job reads 100 TB. How much of that passes through the coordinating process?
    None of it, unless the program asks for records to be returned. The input is divided by description — which range of which file each piece covers — and every byte is read by the worker process assigned that piece. The coordinating process holds the plan and the completion record, both of which grow with the number of pieces rather than with the bytes.
  • The program does some ordinary local work between two distributed steps. Where does that run?
    In the coordinating process, single-threaded on one machine. That is why a loop that pulls intermediate values back and computes over them in plain code does not speed up when the cluster grows: only the distributed steps are spread over the workers, and everything between them is one process's work.
  • Why does the number of pieces affect the coordinating process at all?
    Each piece costs an entry in the plan, an assignment record, a completion report and usually a few control messages. At a few hundred pieces this is invisible; at hundreds of thousands per round of work it becomes real memory and a real share of the round's elapsed time, which is one reason cutting the input ever finer stops paying.

saying these in an interview costs you the question

  • Thinks the coordinating process also runs a share of the data like a worker
  • Believes every output record of the job passes through it
  • Sizes it from the input volume rather than from what is returned
  • Says the job carries on normally once that process is gone
  • Assumes records moving between workers are routed through it
  • Treats it as interchangeable with any other process in the job
open as a page

A job filters a billion rows down to four million and returns them all to the submitting program — what fails first?

level: middleimportance: must knowfreq 66%

basics

~20 s

The coordinating process runs out of memory. Every returned record leaves the workers and lands in that one process's heap on one machine, so the cost follows the result size, not the cluster size — adding worker processes does not help.

open as a page

Should the coordinating process run inside the cluster beside the workers or on the machine that submitted the job?

level: middleimportance: should knowfreq 48%

basics

~20 s

Inside the cluster, the coordinating process sits close to the workers and outlives the submitting machine. Outside it, the session stays interactive but every scheduling message crosses a slower link, and closing the submitting machine ends the job.

open as a page

A long-running job loses the coordinating process on one machine and a worker process on another — how do the two losses differ?

level: seniorimportance: should knowfreq 58%

basics

~20 s

Losing a worker process costs the pieces it was running and the intermediate output it still held; the rest of the job continues. Losing the coordinating process leaves nothing planning, handing out or tracking work, so the job stops.

open as a page