Chapter 6 · Distributed inference
Parallelism, collectives, and topology
6.1

Parallelism, collectives, and topology

Splitting a model across GPUs means picking which axis to cut along. cuts inside a layer: Megatron-LM partitions the weight matrices of a transformer's MLP block by splitting the first GEMM’s weight matrix along its columns and the second GEMM’s along its rows, so each GPU computes a partial result that has to be combined before the next operation can proceed. That combination is a synchronization point after every parallelized block, which is why tensor parallelism wants the fastest possible link between GPUs, typically NVLink inside a node, and degrades quickly across a slower network. cuts the other way: whole layers are assigned to different GPUs, and activations flow between stages instead of being combined mid-layer. Megatron-LM describes its own tensor parallel approach as “orthogonal and complimentary to pipeline model parallelism,” meaning production systems compose both at once rather than choosing one.

Neither form of parallelism does anything without a way for GPUs to exchange the partial results. That is the job of a : an operation like all-reduce, all-gather, or reduce-scatter that every participating GPU executes together. NCCL, NVIDIA’s library for this, implements exactly these primitives: all-reduce, all-gather, reduce, broadcast, reduce-scatter, and arbitrary send/receive patterns, and is “optimized to achieve high bandwidth on platforms using PCIe, NVLink, NVswitch, as well as networking using InfiniBand Verbs or TCP/IP sockets,” supporting any number of GPUs in a single node or spread across many. Tensor parallelism’s per-layer combine step is an all-reduce; pipeline parallelism’s stage-to-stage handoff is a point-to-point send/receive. Which collective a parallelism strategy needs, and how often it needs it, is what determines whether that strategy can tolerate a slow interconnect or requires the fastest one available.

In code, a collective is less exotic than the name suggests. The first complete example in NCCL’s documentation has a single process drive four GPUs: it creates one communicator per device with ncclCommInitAll, then issues an all-reduce on every device inside a group call, which NCCL requires when one thread manages multiple GPUs. Each call names a send buffer, a receive buffer, an element count, a datatype, and a reduction operator; the library sums the four send buffers and leaves the identical result in every receive buffer. Nothing is finished when ncclGroupEnd returns: the operations are queued on CUDA streams, so the host synchronizes each stream before trusting the result.

ncclAllReduce across four GPUs, from the NCCL documentation
ncclComm_t comms[4];
int nDev = 4;
int size = 32*1024*1024;
int devs[4] = { 0, 1, 2, 3 };

//initializing NCCL
NCCLCHECK(ncclCommInitAll(comms, nDev, devs));

//calling NCCL communication API. Group API is required when using
//multiple devices per thread
NCCLCHECK(ncclGroupStart());
for (int i = 0; i < nDev; ++i)
  NCCLCHECK(ncclAllReduce((const void*)sendbuff[i], (void*)recvbuff[i],
      size, ncclFloat, ncclSum, comms[i], s[i]));
NCCLCHECK(ncclGroupEnd());

//synchronizing on CUDA streams to wait for completion of NCCL operation
for (int i = 0; i < nDev; ++i) {
  CUDACHECK(cudaSetDevice(i));
  CUDACHECK(cudaStreamSynchronize(s[i]));
}

//finalizing NCCL
for (int i = 0; i < nDev; ++i)
  ncclCommDestroy(comms[i]);

This is why topology is not a footnote to parallelism, it is a constraint on which parallelism strategies are viable at all. A group of GPUs connected by NVLink can absorb tensor parallelism’s frequent all-reduces; GPUs connected only by a network fabric across nodes generally cannot, which is why tensor parallelism is typically confined inside a node and pipeline or data parallelism is used to scale beyond it, where the coarser, less frequent communication tolerates the added latency.

Tensor vs. pipeline parallelism. The same four-layer model split two ways. Tensor parallelism cuts every layer in half across both GPUs, so activations must be all-reduced at every layer boundary. Pipeline parallelism gives each GPU whole layers and sends activations across the boundary once.Illustrative numbers

The size of that fastest tier is a moving target. NVIDIA’s multi-node tuning guide notes that before the GB200 NVL72, an NVLink domain topped out at eight GPUs on an HGX H200 baseboard at 900 GB/s of communication per GPU; the NVL72 rack design extends one domain to 72 Blackwell GPUs at 1.8 TB/s each, with the NVLink switch providing 130 TB/s of aggregate GPU bandwidth inside the domain, and describes that 72-GPU domain as acting as one massive GPU. For the parallelism strategies above, that is a ninefold change in where the scale-up boundary sits: a split that would have crossed the network fabric between eight-GPU nodes can now stay on NVLink across a whole rack.

NVLink is one vendor’s fabric, and open standards now exist for both directions of scaling. UALink is the scale-up one: its 200G 1.0 specification “defines a low-latency, high-bandwidth interconnect for communication between accelerators and switches in AI computing pods” and “enables 200G per lane scale-up connection for up to 1,024 accelerators within an AI computing pod.” Ultra Ethernet is the scale-out counterpart, a Linux Foundation project whose stated mission is an “Ethernet based open, interoperable, high performance, full-communications stack architecture” for AI and HPC at scale; the problems it names, multi-pathing, fast reaction to congestion, and flows “where tail latency is the figure of merit,” are this chapter’s constraints restated as networking requirements. Whichever fabric wins a given deployment, the collective at the top of the stack stays the same; what the standards change is who can build the hardware underneath it.

Scale-up and scale-out are not informal labels, either. The Ultra Ethernet specification defines them, and the definitions are the vocabulary this section has so far been gesturing at. It differentiates three network types: a frontend network, a backend scale-out network, and a scale-up network. The frontend network is the operational network in datacenters that “connects all compute nodes to the outside world,” other datacenters or customers on the internet. The backend is “a specialized high-performance network of limited scope relative to the frontend network,” usually deployed across a cluster, usually its own layer-3 subnet, and usually not connected to the frontend directly. Scale-up networks are “typically very specialized short-range interconnects that often come with only a single tier of switches or possibly no switch at all,” and the spec names NVLink alongside AMD’s XGMI, Intel’s Xe Link, switched PCI Express, and CXL as its modern examples. Its table of 2024 characteristics separates the three by one-way latency requirement, 100 microseconds and up for the frontend, under 10 microseconds for backend scale-out, and under 1 microsecond for scale-up, and by maximum link length, 1500 m, 150 m, and 5 to 10 m. Ultra Ethernet then states which one it is for: the specification “is focused on the backend scale-out network,” with support for the other two considered only opportunistically. That is the same boundary this section has been drawing from the other side. Tensor parallelism wants the scale-up network; what crosses a node gets the scale-out one.

Turning the split on#

In a serving engine both cuts are flags. vLLM’s guidance for a single model replica is a ladder rather than a menu: if the model fits on one GPU, distributed inference is probably unnecessary; if it does not fit on one GPU but does fit on one node, use tensor parallelism and set tensor_parallel_size to the number of GPUs in that node; and only when the model is too large for a single node does the guidance reach for both axes at once, with tensor_parallel_size set to the GPUs per node and pipeline_parallel_size to the number of nodes. Composing the two is a second flag, not a second program.

tensor and pipeline parallelism composed on eight GPUs, from vLLM's parallelism and scaling guide
# Eight GPUs total
vllm serve gpt2 \
     --tensor-parallel-size 4 \
     --pipeline-parallel-size 2

The command says nothing about all-reduces, but the two numbers decide which collectives run and how often. Four-way tensor parallelism puts an all-reduce after every parallelized block inside a stage; two-way pipeline parallelism adds one activation handoff between the two stages. The same pair of flags crosses a machine boundary as well, which is why the multi-node recipe sets tensor parallel size to the GPUs in a node and pipeline parallel size to the node count: the frequent collective stays on the fast fabric, and the infrequent one takes the network.

The guidance has two edge cases, and both are about what the hardware will bear rather than what the model needs. If the model fits inside a node but the GPU count does not divide it evenly, vLLM recommends pipeline parallelism, which splits along layers and supports uneven splits, with tensor_parallel_size set to 1 and pipeline_parallel_size set to the number of GPUs. And if the GPUs in a node have no NVLink between them, the L40S being the example given, the same advice arrives for a different reason: pipeline parallelism yields higher throughput and lower communication overhead there. Tensor parallelism is not the within-node default because it is inherently better. It is the default because a node usually has NVLink.

Whether enough GPUs were provisioned is not a calculation the operator has to do alone, because the engine reports it. Once the server is up, vLLM logs the total number of tokens its GPU KV cache holds and an estimate of how many requests that supports concurrently at the model’s configured maximum sequence length: the documentation’s own example reads a GPU KV cache size of 643,232 tokens and a maximum concurrency of 15.70x for 40,960 tokens per request. Those two lines answer the question the flags were chosen to ask. If they come in under what the deployment needs, the remedy the guide gives is not a cleverer split but more hardware: add GPUs or nodes, then set the two sizes again.

What NCCL decides for you#

Underneath those flags sits a library that has already made most of the decisions. NCCL sorts its environment variables into two categories, and the split is a statement of intent: some are needed to make NCCL follow system-specific configuration and can be kept in scripts and system configuration, while the ones under its debugging heading “should not be used in production nor retained in scripts, or only as workaround, and removed as soon as the issue is resolved.” Keeping those set, the documentation warns, may result in sub-optimal behavior, crashes, or hangs.

NCCL_P2P_LEVEL sits on the debugging side, and reading its accepted values is the fastest way to see the interconnect hierarchy the library actually works in. The level is a maximum distance between GPUs beyond which peer-to-peer transport is not used, and the distances have names: LOC never, NVL when GPUs are connected through NVLink, PIX when they share a PCI switch, PXB across potentially several PCI switches, PHB when they sit on the same NUMA node so traffic goes through the CPU, and SYS between NUMA nodes, potentially crossing the SMP interconnect. Left unset, NCCL “will attempt to optimally select a value based on the architecture and environment it’s run in.” That ladder is the same one this section has been describing, expressed as a cutoff.

What does belong in a production script is the part NCCL cannot infer: which interfaces it is allowed to use. NCCL_SOCKET_IFNAME filters IP interfaces by prefix, with ^ excluding and = forcing an exact name, and its automatic selection already favors interfaces starting with ib over the rest. NCCL_IB_HCA does the same for RDMA adapters, where the documentation recommends always adding the = prefix so that a token like mlx5_1 does not also select mlx5_10 through mlx5_19, and notes a fixed upper limit of 32 such devices. When the topology itself is in question, NCCL_TOPO_DUMP_FILE writes out the XML NCCL detected, which on a multi-node NVLink system contains the full NVLink domain. The variables worth keeping describe the machine; the ones that override NCCL’s own choices are how a bug gets diagnosed, not how a cluster gets tuned.