Two worker processes of one job sit on the same machine — does moving records between them still cross a network path?
answer
- the boundary is the process
- same machine, different address space
- wire skipped, encoding still paid
- only slots in one process pass references
basics
~20 sMostly yes. Two worker processes have separate address spaces, so records are encoded into bytes at one end and decoded at the other even on one machine. Sharing a machine skips the wire, not the encoding — only slots inside one process avoid it.
solid answer
~50 sThe boundary that costs money is the **process**, not the machine. Two worker processes — the processes that run pieces of the job and own the memory those pieces use — have separate address spaces, so a record going from one to the other must be turned into bytes and turned back. Co-location removes the wire between machines, so bandwidth is better and latency lower, but the encoding, the buffer copies and the per-connection bookkeeping are paid identically. What the transport looks like varies: some engines have the consumer read a local file the producer wrote, some stream records over a connection whether or not the ends share a host, and some detect co-location and take a shorter path. Only two work slots **inside one process** can hand a record over as a reference, with no encoding at all — which is the concrete reason more slots per process reduces exchange cost.
go deeper
Remember that separate processes do not share memory, so records passing between two worker processes are turned into bytes and back even on one machine.
Separate the two boundaries cleanly: the process boundary costs encoding and decoding, the machine boundary adds the wire on top of it.
Resist designing around co-location you cannot control, and argue instead from the shape of the workers, which changes the mix of handoffs deterministically.
Weigh whether the platform should concentrate slots per process fleet-wide, given what the organisation's jobs actually spend on encoding versus what a larger blast radius costs it.
## The boundary that matters is the process When a step needs records held by other workers, the job has to redistribute them. It is tempting to think of that cost as a network cost and therefore to think that two worker processes on one machine exchange records almost for free. They do not, and the reason is that the expensive boundary is not the machine — it is the **operating-system process**. A worker process owns its own address space. A record living inside it is a structure at an address that means nothing to any other process. To hand it to a different process, the producer must turn it into a flat sequence of bytes, and the consumer must turn those bytes back into a record it can operate on. That work happens whether the two processes are in different data halls or on the same host. ## What sharing a machine actually saves - **The wire.** No switch, no cross-machine hop, so the bandwidth available is much higher and the latency much lower than a remote fetch. - **Exposure to the network.** A same-machine transfer does not compete with the rest of the cluster's traffic and does not fail for a reason out in the fabric. - **Sometimes a copy or two**, where the runtime recognises co-location and uses a local file read or a shared memory region rather than a loopback connection. ## What it does not save - **Encoding and decoding each record.** On many workloads this, and not the transfer, dominates the cost of moving data. - **Buffering and bookkeeping.** The producer still assembles output destined for that consumer; the consumer still tracks what it has fetched from whom. - **The barrier behaviour of the step itself.** Whether the consumer can start before the producer has finished is a property of the step, not of where the two ended up. | Path | Wire | Encode/decode | Handoff | |---|---|---|---| | Two processes, different machines | yes | yes | bytes over a connection | | Two processes, same machine | no | yes | loopback, or a local file read | | Two slots, same process | no | no | a reference in memory | ## Where designs differ This is exactly the place where one engine's behaviour is routinely mistaken for the class's: - Some runtimes **materialise** intermediate output to local files and have the consumer fetch it afterwards. When both ends sit on one machine, that fetch can be a plain local read, and in some implementations a mapped file, which is genuinely cheaper than a loopback connection. - Others **push records across a connection as they are produced**, with nothing written down in between. Those implementations often use the same code path regardless of placement, so co-location saves the wire and nothing else. - Where a **detached intermediate-output server** is in use — a separate long-lived process on the machine that keeps serving a worker's output after that worker is gone — even a same-machine fetch goes through that third process rather than directly. - Whether a short-circuit path for a co-located consumer exists at all is an implementation choice, and where it exists it is often limited to particular operators. So the honest general claim is narrow and reliable: **crossing a process boundary always costs the encode and the decode; crossing a machine boundary additionally costs the wire.** ## Why this matters for shaping a worker The practical conclusion is not "try to co-locate". A job does not choose where its processes land — whatever grants the machines does — and even when two of them share a host, the job cannot depend on it from run to run, nor know which pairs of pieces will need to talk. The lever the author does have is the **shape**. Putting more work slots inside one process converts some fraction of the job's handoffs from cross-process transfers into in-memory references, and that fraction is reliable, not accidental. It is one of the real arguments for the few-and-large shape, and it should be stated that way — as an effect of concentration, not as a hope about placement. ## The wrong version of this answer Two claims fail an interview here. The first is that co-located processes share memory and so exchange records for free: they do not, because separate processes do not share an address space. The second is that co-location makes no difference at all: it does, because the wire is real and the loopback path is much faster. The correct answer names the two boundaries separately and says which cost sits at each.
- Can a job arrange for two of its worker processes to land on the same machine to make exchanges cheaper?Generally not. Placement belongs to whatever grants the machines, and even where a hint exists the job cannot know in advance which pairs of pieces will need to exchange records. Treat co-location as luck rather than as design. The reliable version of the same idea is to put more work slots inside one process, which turns some handoffs into in-memory references by construction.
- If encoding is the dominant cost, why does moving a step's data between machines still hurt so much?Because the wire adds a shared, finite resource to a cost that was already there. Every worker sending to every other means the traffic scales with the square of the width, all of it competing for the same links, and a slow or saturated path stalls the consumers waiting on it. The encoding cost is per record and local; the wire cost is per record and contended.
saying these in an interview costs you the question
- Says two processes on one machine exchange records straight from shared memory
- Claims co-location removes the encoding and decoding cost of an exchange
- Assumes a job can place its worker processes together on purpose
- Thinks a same-machine handoff costs exactly what a cross-machine one costs
- Believes every engine uses the same transport regardless of where the consumer sits