skip to content

How do you split one large LLM across two 8-GPU nodes for serving?

level: seniorimportance: should knowfreq 42%

answer

  1. match the strategy to the link
  2. heavy collectives stay inside the box
  3. one tensor crosses the network
  4. ask first whether it must span nodes
  5. pipelines need traffic to stay full

basics

~20 s

Keep tensor parallelism inside each node — TP=8 over the intra-node fabric — and cross the network with pipeline parallelism, PP=2. The node boundary is the slow link, and only pipeline parallelism sends little enough traffic to survive it. First check the model cannot simply fit on one node.

solid answer

~60 s

The layout follows the topology. Inside a node the GPUs share a high-bandwidth fabric, so that is where the per-layer all-reduces of tensor parallelism belong: TP=8. Between nodes you have a network link that is one to two orders of magnitude slower and higher latency, and tensor parallelism would put roughly two collectives per layer across it for every token — so instead you place a pipeline boundary there: PP=2, sending one activation tensor per stage per micro-batch. That is the standard 16-GPU layout. Three checks before committing. **Can you avoid the split entirely?** A quantized checkpoint that fits in 8 GPUs and runs as two independent replicas will almost always beat a 16-GPU single instance on throughput per dollar and on blast radius. **Does your traffic have the concurrency to fill a pipeline?** Stages idle without it, so a low-QPS latency-sensitive service gets nothing from PP. **Is the fabric actually fast?** Cross-node tensor parallelism is only worth discussing on a high-end RDMA fabric, and even there pipeline parallelism across the boundary is usually the safer choice.

go deeper

for a junior

Know that GPUs inside one machine talk to each other far faster than machines talk over a network, and that how you split a model has to respect that difference.

for a middle

Explain why tensor parallelism stays inside a node — a collective per layer per token — while pipeline parallelism, which sends one activation per stage, is what crosses the network.

for a senior

Give the TP=8 x PP=2 layout and defend it, then raise the operational consequences: gang scheduling, minutes-long cold start, a fused failure domain, and the concurrency a pipeline needs to pay off.

for a principal

Own whether the multi-node instance should exist: weigh quantizing to fit one node and running replicas against the fabric investment, blast radius and scheduling complexity a 16-GPU serving unit imposes on the platform.

## Read the topology first Multi-node serving is a placement problem, and the placement is dictated by where bandwidth drops. Inside a modern GPU node, peer bandwidth is in the hundreds of GB/s with microsecond-scale latency. Between nodes, even a good fabric is far behind that, and a commodity Ethernet fabric is not close. Every layout decision follows from matching each parallelism strategy to the link it can tolerate. ## The default answer Tensor parallelism generates roughly two all-reduces per transformer layer per token — the heaviest, most latency-sensitive traffic in inference. It goes inside the node: TP = 8, the full node. Pipeline parallelism sends one activation tensor per stage boundary — a slab of size (tokens x hidden_size x dtype bytes), often well under a megabyte, point-to-point, and not on a per-layer cadence. It goes across the node boundary: PP = 2, node 0 running the first half of the layers, node 1 the second half. So 16 GPUs become TP=8 x PP=2. This is the layout to state first, because it is right in the overwhelming majority of cases. ## Sanity checks before you build it **1. Do you need 16 GPUs at all?** A 16-GPU single instance is the most expensive, most fragile serving unit you can construct: one failed GPU anywhere kills it, deploys are all-or-nothing, and scheduling requires two whole nodes simultaneously. If quantizing the weights lets the model fit in one node, two independent 8-GPU replicas give you more throughput, half the blast radius and vastly simpler operations. The senior answer always considers not sharding across nodes. **2. Do you have the concurrency?** Pipeline parallelism only pays when enough requests are in flight to keep both stages busy; otherwise one node is idle while the other works, and single-request latency is slightly *worse* than it would be on a hypothetical single node. High-throughput batch workloads fill a pipeline happily. A low-QPS interactive endpoint does not. **3. What is the fabric?** Cross-node tensor parallelism is not categorically impossible — with a high-end RDMA fabric and GPUDirect it can be made to work — but the collectives are being asked to cross a link that is much slower than the one they were designed for, and the result is usually worse than the pipeline layout on the same hardware. On Ethernet without RDMA it is not a serious option. ## Operational consequences to raise unprompted - **Failure domain.** All 16 ranks form one process group. A hung rank presents as a stalled collective, not a clean crash, so you need liveness checks with timeouts rather than relying on the process to exit. - **Startup.** Two nodes must both pull tens of gigabytes of weights and rendezvous before the server is ready. Cold start is measured in minutes, which interacts badly with aggressive autoscaling. - **Scheduling.** You now need gang scheduling: sixteen GPUs across two specific nodes, available at once. Partial placement is a deadlock waiting to happen. - **Network isolation.** Noisy neighbours on a shared fabric become your latency problem, and inter-node traffic should be on a dedicated high-speed interface, not the general-purpose one. ## How the engines handle it Multi-node support differs sharply and is worth naming rather than assuming. vLLM supports multi-node serving with a distributed executor backend — Ray or multiprocessing — coordinating ranks across machines, and exposes the two degrees as `--tensor-parallel-size` and `--pipeline-parallel-size`. TensorRT-LLM bakes the parallel layout into the compiled engine, so the split is decided at build time and changing it means rebuilding. Whichever stack you use, the layout must be identical on every node and the weights must be reachable from both. ## What a weak answer looks like "Set the tensor-parallel size to 16." It is a legal configuration and it will start, which is what makes it a good trap — it just runs badly, because you have placed the per-layer collectives on the slowest link in the system. The strong answer names the boundary, matches the strategy to the link, and then argues about whether the split was necessary at all.

  • What would make you reconsider spanning nodes at all?
    If quantizing the weights brings the model within one node's memory with usable KV-cache headroom, run two independent 8-GPU replicas instead. That doubles the number of failure domains rather than fusing them, removes cross-node traffic from the serving path entirely, simplifies scheduling and deploys, and usually delivers more aggregate throughput per dollar. Validate the quantized checkpoint on task metrics first, since the whole argument rests on quality holding up.
  • How does cold start change when a serving unit spans two nodes?
    Both nodes must load tens of gigabytes of weights and then rendezvous before any request is served, so readiness is gated by the slower node and typically takes minutes. That makes scale-to-zero impractical and forces autoscaling to react on leading indicators well before saturation. Keeping warm capacity, pre-pulling images and staging weights on fast local storage are the standard mitigations.
  • You must span nodes and the fabric is plain Ethernet without RDMA. What do you accept?
    Pipeline parallelism across the boundary and nothing heavier — cross-node tensor parallelism on that fabric will be dominated by collective latency. Accept that the pipeline needs sustained concurrency to keep both stages busy, so single-request latency will not improve and may worsen slightly. Size micro-batching to keep stages full, put inter-node traffic on a dedicated interface, and monitor for stage stalls that indicate the link is the bottleneck.

saying these in an interview costs you the question

  • Setting the tensor-parallel size to span both nodes
  • Assuming pipeline parallelism improves single-request latency
  • Ignoring that all sixteen ranks form one failure domain
  • Planning scale-to-zero for a multi-node serving unit
  • Not asking whether the model could fit in a single node

context