TL;DR
- NCCLX is the communication library that sits under PyTorch in Meta's Llama 4 training and inference. It is built on NCCL and adds a second transport, CTran, in which a CPU thread runs each collective and the network card moves data directly between the application's GPU buffers.
- Baseline NCCL runs collectives inside a GPU kernel and copies every message through its own buffers. Across several data-center buildings that design takes GPU threads away from the model, limits how much data one network request can carry, and puts control messages from the receiver on the critical path of a transfer.
- On CTran the paper builds one custom collective for each kind of parallelism: zero-copy send and receive for pipeline parallelism, a window-and-put pipeline that hides tensor-parallel transfers inside the matrix multiply (1.57ร lower latency on one node), a fault-tolerant AllReduce for data parallelism, and an AllToAll for mixture-of-experts inference that reads its sizes on the GPU (decode time down by up to 83%).
- A large part of the paper is about operating at this size. Creating the default process group for 96K GPUs takes 24 s instead of 265 s, the library's GPU memory use is cut by almost half, a fault analyzer finds the first collective that stalled, and 96K ranks are tested on CPU servers.
A 100k-GPU job is a network problem
A training cluster of more than 100,000 GPUs does not fit in one data-center building. Meta joins several nearby buildings into a single RoCE fabric (Section 2.3). Inside a building the network is a three-layer Clos: a Rack Training Switch (RTSW) connects the GPUs of one rack, Cluster Training Switches (CTSW) connect the racks of one AI zone, and Aggregator Training Switches (ATSW) connect the zones of the building. Between buildings, the ATSW layers form a fully connected mesh.
The further apart two GPUs sit, the slower the link between them. Against two GPUs in the same rack, the paper gives the latency as 7ร across racks in one zone, 15ร across zones and 30ร across buildings, from the switching delays and cable lengths that add up along the path. Traffic that leaves a zone also shares its uplinks: the cross-zone over-subscription ratio is 1:2.8, down from 1:7 in the network used for Llama 3, and traffic between buildings sees the same ratio.
Pick where the second GPU is. The path shows the switch layers the data crosses, and the bars are the paper's latency ratios, drawn to scale.
Large models are trained with several kinds of parallelism at once, and each kind produces different traffic (Section 2.1; the kinds are compared in data vs tensor vs pipeline parallelism). Collectives of the innermost domains, such as tensor parallelism, are fully exposed in the training step and stay on high-bandwidth links. The middle domains, expert and pipeline parallelism, send medium-sized messages that are partly hidden behind computation. The outermost domains, fully sharded and plain data parallelism, move the most data but are usually well hidden by the work of the inner domains. Meta's job scheduler is aware of the topology: it places consecutive ranks as close together as it can, and a job can state how many GPUs it wants at each level, so that each kind of parallelism lands on a chosen level of the network (Section 5.4).
The spread of latencies sets a requirement that runs through the rest of the paper. To use a long link well, the library has to keep enough data in flight to fill the link's bandwidth-delay product, and not so much more that the switches along the path congest.
What NCCL does, and where it stops
NCCL is host-initiated: the CPU schedules a collective, and its arguments (sizes, data types) are CPU variables. It was designed for the bulk, regular collectives of classic distributed training, AllReduce, AllGather and ReduceScatter, and it favors bandwidth over latency (Section 2.2). The paper names two properties of its design that held back training at the scale of Llama 4 (Section 4):
- It is kernel-driven. A collective's algorithm runs mostly inside a CUDA kernel, and a proxy thread on the host issues the network operations for it.
- It is copy-based. Every message is copied into a buffer that NCCL owns before it is sent, and out of one after it is received.
A third limit matters for inference. The arguments of a collective must be known on the host when the collective is enqueued. If a size is computed by an earlier GPU kernel, it has to be copied to host memory first. That costs a CPU-GPU synchronization, and it does not work inside a CUDA graph, where the usual workaround is to pad every message to its largest possible size.
NVIDIA's other library, NVSHMEM, has the opposite model: communication is started from inside a GPU kernel, with arguments that live on the GPU. That gives low latency and dynamic sizes, but it needs a symmetric memory region allocated and registered on every rank, and that memory is not shared with PyTorch. Applications with mixed traffic end up running both libraries. The paper counts the costs of that: the two runtimes cannot be optimized together, they cannot share resources, and operating two of them doubles the maintenance.
NCCLX is meant to be the single stack for both. It offers three kinds of API, each with collectives, point-to-point and remote memory access: host-initiated calls, host-initiated calls whose metadata lives on the GPU, and device-initiated calls. The third kind is ongoing work and is not described. Under the API a call goes either to baseline NCCL or to CTran. Custom operations exist only in CTran, and for the classic collectives the choice is made by environment variables, usually set by offline tuning before a model launches.
CTran: the CPU drives, the GPU does not copy
Host-driven collectives. CTran starts one background CPU thread for each communicator. When the program calls a collective, CTran runs the algorithm on that thread and launches a small "stall kernel" on the caller's CUDA stream, so that the collective keeps its place in the order of the stream. The two sides signal each other through a flag in host-pinned memory at the start and at the end (Figure 3). There are three cases:
- A collective that only uses the network, such as the AllGather between nodes in fully sharded data parallelism or the send and receive of pipeline parallelism, is run entirely by the CPU thread. Each RDMA operation is posted with no coordination with the kernel.
- A collective that uses the network and NVLink with no dependency between the two, such as the AllToAll of a mixture-of-experts layer, lets the kernel copy over NVLink while the CPU thread posts RDMA operations in parallel.
- A collective whose steps depend on each other, such as an AllReduce between nodes where the kernel's reduction and the network transfer form a pipeline, synchronizes the kernel and the CPU thread through producer-consumer flags. The paper measures that overhead as always under a microsecond.
Running the algorithm on the CPU also makes the collective algorithms of classic HPC easy to bring over. Before NCCL gained the PAT algorithm in release 2.23, it had only a ring for AllGather and ReduceScatter. The authors ported latency-oriented algorithms from the MPI literature instead: Bruck and recursive doubling for AllGather, recursive vector-halving with distance doubling for ReduceScatter, and a tree for Broadcast (Section 4.3.2). For the recursive-doubling AllGather they found that contacting the farthest peers first works best on their over-subscribed network.
Zero-copy transfers. In baseline NCCL a message makes four hops around the network transfer (Figure 4a). The sender's GPU copies it from the application's buffer into a FIFO buffer that NCCL registered with the network. The network card reads that buffer over PCIe and sends the data. The receiving card writes it into the receiver's FIFO buffer, and the receiver's GPU copies it into the application's buffer. The two copies inside the GPUs run on the streaming multiprocessors (SMs) of the collective's kernel and use the memory bandwidth of the GPU. The PCIe hops are done by the network card and need no GPU threads.
The paper lists three costs of the copies:
- They compete with the model for GPU threads and memory bandwidth.
- To overlap the fast copies with the slower network, the message is cut into chunks, and each RDMA operation carries one chunk. Every stage of that pipeline needs a GPU-CPU synchronization, and the small requests cannot fill a long, high-latency link.
- Each pipeline, called a channel, is bound to one GPU thread block and has its own FIFO buffers. Using more thread blocks to speed up a collective therefore costs more GPU memory.
CTran sends the data from the application's buffer on one GPU to the application's buffer on the other (Figure 4b). A collective that only uses the network needs no work from a GPU kernel at all. The copy-based path also depends on clear-to-send messages from the receiver while a transfer runs, and on a long link each of them costs the latency of that link. A zero-copy transfer needs a single control message from the receiver (Section 4.4).
Play both transfers and compare what moves the data at each hop.
In the paper's point-to-point benchmark the extra copies add a constant amount to every transfer, enough to double the transfer time between two hosts for small messages (Section 4.5). For medium-sized messages the copy-based path does not reach a reasonable share of the network's bandwidth, even after NCCL's chunk size and channel count are tuned for the benchmark. The authors add that such tuning is not usable in production, because the settings are global and can slow down other collectives in the same model.
Registering the buffers. Zero-copy is routine in CPU-based HPC. On GPUs under PyTorch it has a cost of its own: before a network card can read or write a buffer directly, the buffer's physical pages have to be looked up and pinned. The paper puts that at hundreds of microseconds to a few milliseconds for the buffer sizes common in LLM workloads, with occasional spikes to 100 ms that it traces to lock contention inside the GPU RDMA driver. The answer is an extension of PyTorch's CUDA caching allocator with two modes (Section 4.3.1):
- Automatic registration. CTran records every memory segment the allocator creates and registers a tensor with the network the first time a collective uses it. This suits stages that reuse the same memory, such as the AllGather at the start of every step in fully sharded data parallelism.
- A memory pool. A large pool is allocated and registered once, and tensors that will be communicated are taken from it. This suits pipeline parallelism, where memory pressure makes the allocator remap segments often and automatic registration would register again and again. Tensors have to be labelled for the pool, and when the default pool is close to running out of memory it can borrow the free space of the registered pool.
Zero-copy needs its own flow control
Handing a whole message to the network at once creates a problem that the copy-based path did not have. Meta's fabric does not use a conventional congestion control scheme such as DCQCN. It relies on switches with deep buffers to absorb bursts and on flow control driven by the receiver, inside the collective library, to prevent lasting congestion (Section 4.4). The chunking of the copy-based path gave that flow control for free. The paper reports that zero-copy alone gave suboptimal performance, and attributes it to large messages building up in the switch buffers.
The answer is Dynamic Queue Pair Load Balancing (DQPLB). The transfer stays zero-copy, but CTran cuts the message into segments and limits how many are in flight:
- A connection has one control queue pair, which exchanges memory addresses at the start of a collective, and one or more data queue pairs.
- Three limits apply per connection: the number of data queue pairs, the number of unacknowledged segments on each of them, and the size of a segment.
- The limits are set per distance, for the four distances of the figure above. Close connections get conservative settings, because their bandwidth-delay product is small. Distant ones get more queue pairs and more segments in flight.
The sender hands segments to the data queue pairs in turn. When a queue pair reaches its limit, nothing more is posted on it until one of its segments completes. In the released code the sender skips a full queue pair and moves on to the next, so a queue pair on a less contended path ends up carrying more of the message.
Because the segments travel on different queue pairs, they can arrive out of order. Each one is sent as an RDMA write with 32 bits of immediate data: bits 0 to 23 hold a sequence number, bit 30 marks the fast path, and bit 31 marks the last segment of a message. The receiver keeps the next sequence number it expects. A segment that arrives early is stored in a hash map, and when the expected one arrives it is processed together with every consecutive segment already waiting. The completion of a message is therefore reported only after all of its segments are in.
The figure plays one message of 16 segments over four queue pairs, one of which is slow.
Small, frequent messages take a fast path instead. They are sent on a single dedicated queue pair, which keeps them in order, so the receiver only has to count them. Together with load-balancing work from an earlier paper and tuning of the spine switches, the authors report that the buffer build-up in the switches fell by an order of magnitude from Llama 3 training to Llama 4 training.
Training: one collective per kind of parallelism
Pipeline parallelism: send and receive without the GPU
Pipeline parallelism passes activations between stages with point-to-point sends, often across racks and sometimes across zones, and each receive is followed by computation on the same GPU (Section 5.1). With baseline NCCL the authors saw two problems. The copy-based path cut the data into chunks of 128 KB to 512 KB, too small to hide a latency 7 to 15 times that of a rack. The staging copies also occupied 4 thread blocks of 640 threads each for a send or receive that only used the network, which slowed the matrix multiplies running beside it.
In CTran the receiver's CPU thread sends the registration of its receive buffer to the sender, and the sender posts one RDMA write from its buffer to the receiver's. No GPU thread is involved, and the network sees the whole message, so sends of tens of megabytes reach the peak bandwidth of the network. For messages of 1 MB to 128 MB the paper measures a speedup of 1.09ร to 2.7ร over the copy-based send (Figure 10). NCCL's own zero-copy path performs about the same in the benchmark, but the authors could not use it in production because of how it registers buffers under the allocator's expandable segments.
Tensor parallelism: hide the transfer inside the multiply
Llama training uses tensor parallelism in the style of Megatron-LM: the input of a pair of matrix multiplies is gathered from all GPUs of the group with an AllGather, and the output is combined with a ReduceScatter. These collectives sit directly on the critical path of every layer. Earlier ways of overlapping them with the multiply either work only inside one host, because they depend on CUDA inter-process communication, or start the transfers from inside modified multiply kernels, which are slower than the vendor's tuned library (Section 5.2).
CTran adds two pieces with the semantics of MPI's one-sided communication:
- A window: every rank registers a memory region of the same size in advance and gives its address and access key to the other ranks.
- A put: any rank writes into any peer's window. A put uses NVLink's copy engine inside a node and RDMA between nodes, and neither uses GPU threads.
With a window in place, the gathered input can arrive in chunks, and the multiply can start on each chunk as soon as it lands. The paper arranges the chunks as a tree that follows the hardware. With two nodes of four GPUs, a rank first exchanges one chunk with a neighbor in its node and multiplies it, then receives two chunks and multiplies a tensor twice as large. Meanwhile the first chunk from the other node is already on its way. A transfer between nodes over RDMA is about 8 times slower than one over NVLink, so that chunk arrives when the work inside the node is done, and its wait is hidden.
On one node, with 32 MB tensors transferred seven times on each rank, the transfer takes about as long with overlap as without, and the multiply is somewhat slower because it works on smaller pieces. End to end the overlapped version has 1.57ร lower latency (Figure 11).
Data parallelism: lose a group, keep training
At 100,000 devices hardware faults are frequent, and what matters is the share of time spent on useful training. Fully sharded data parallelism spreads every parameter and every piece of optimizer state over all workers, with no redundancy, so one failure stops the whole group (Section 5.3).
Together with the model researchers, the authors changed the outermost dimension to hybrid sharding. The GPUs are divided into replica groups. Inside a group, parameters and optimizer state are sharded as before. Across groups, each group trains a full copy of the model on its own inputs, and the groups average their gradients with one AllReduce at the end of each step. When a machine fails, only its group shuts down and the others continue with a smaller ring. When the machines are replaced, a new group of 4K GPUs is formed and joins again. A global coordinator, which talks to the lead of each group over a separate channel, detects faults and manages the membership.
The collective that makes this work is the Fault Tolerant AllReduce (FTAR). Timeouts and error handling are run from CTran's CPU thread. FTAR travels across zones and buildings, where over-subscription is highest, so it uses a ring: each GPU talks to two neighbors only, which keeps the traffic in flight low. The copies and reductions inside the GPU are arranged as a pipeline hidden behind the network transfers, with a fixed chunk size and a fixed number of chunks, so that the amount of data in flight between two peers has a known maximum. The paper settles on 8 MB chunks and 2 thread blocks of 512 threads. Against NCCL's AllReduce, which uses 4 thread blocks of 544 threads, FTAR has comparable latency with half the thread blocks, and 9% to 18% lower latency when NCCL is held to the same two (Figure 12).
Inference: small messages, unknown sizes
Inference moves far less data than training, but it has to answer quickly. The collective that matters is the AllToAll of a mixture-of-experts layer. On each GPU a router kernel chooses the experts for each token, the tokens are sorted by expert into a send buffer, and an AllToAllv sends each one to the GPU that holds its expert (Figure 13).
How many tokens go to each GPU depends on what the router decides, and NCCL copies the send counts on the host when the collective is enqueued, before the router has run. In eager mode the program can wait for the router and copy the counts to the host, at the cost of a CPU-GPU synchronization. Inside a captured CUDA graph, which inference stacks use to avoid kernel launch overhead, that is not possible. The collective then has to send the largest possible count to every GPU, sized for the case where every token goes to the same expert, and most of what it sends is padding.
NCCLX's answer is a collective whose metadata lives on the GPU (Section 6.1). AllToAllvDynamic takes its counts by reference, so they can change until the moment the collective starts to execute, and it returns the counts it received to the caller. Press "Route again": the real sizes change with the router's choice, and the padded sizes do not.
With the padding gone the messages are small, and the cost moves to the CPU. For an AllToAll over RDMA with N ranks, the paper models the latency as
where Tc is the time the CPU needs to prepare a message for one peer, S is the message size per rank and BW is the bandwidth between two ranks. The transfers to different peers overlap, but one CPU thread prepares them one after another, so with many ranks and small messages the first term dominates. A profile of a small-message AllToAll on CTran shows where the time goes (Table 2):
| Step | Share of the latency |
|---|---|
| Exchange control messages with every peer | 50% |
| Issue the RDMA puts | 20% |
| Wait for the puts to be received and completed | 30% |
Libraries built on NVSHMEM, such as DeepSeek's DeepEP, prepare the messages for different peers on parallel GPU threads. The paper keeps the host-driven design and shortens the preparation instead (Section 6.2):
- The critical path. Functions are inlined, error checks become optional in a low-latency mode, and the mode is passed down as a template parameter so that it adds no branches.
- The control messages. A CUDA graph requires its tensors to stay the same between capture and replay, so the memory handles are exchanged once, at capture. The exchange also served as a barrier that told the sender the receive buffer was free. That role is replaced by double buffering in the model code: two consecutive AllToAlls never write to the same receive buffer.
- The puts. Small messages take the fast path of the previous section. Several buffers that are not contiguous are chained into one request, so the lock is taken and the network card is notified once.
The end-to-end test measures decode time with the collective built into the token shuffling code, on 4, 8 and 16 hosts, with one or four experts per token and a batch of 128 or 256 (Table 3). The baseline already runs on CTran and uses two AllGathers and one AllToAll.
The baseline's decode time grows with the number of hosts, and the time with AllToAllvDynamic changes much less, so the gain is largest at 16 hosts. With four experts per token and a batch of 256, 164.29 ms falls to 27.43 ms.
Running it at 100k
Starting the job
Before a job can communicate, every rank has to find the others and exchange state. At fewer than a thousand GPUs that takes some tens of seconds and nobody minds. The work of coordination grows quadratically with the number of ranks, and with baseline NCCL it takes more than four minutes at 96K GPUs. Jobs of this size restart often, so that time is paid again and again (Section 7.1). The paper names the causes:
- Peer discovery is serialized through a bootstrap server. At 100K ranks the last rank waits 100 seconds just to connect to it.
- TCP listen queues overflow beyond 64K ranks, and connections are reset silently.
- Computing the topology is quadratic in the number of ranks: 10 s at 48K ranks, and a projected 100 s at 100K. Building the ring is quadratic as well.
The fixes are spread over PyTorch, baseline NCCL and CTran:
- One global communicator is created, and the roughly ten process groups of a job are derived from it by splitting, in place of bootstrapping each one separately.
- Peer discovery moves from NCCL's bootstrap server to TCPStore, the store PyTorch's process groups already use. At 16K GPUs that cuts topology formation from 18.45 s to 4.1 s.
- The AllGather that exchanges control data during bootstrap runs in both directions, which cuts its steps from N โ 1 to N / 2, and seven such calls are combined into four.
- The quadratic algorithms are rewritten to be linear.
- CTran opens a connection only when a collective first needs it.
For the AllGather change and the linear-time algorithms the paper cites two pull requests against NVIDIA's NCCL repository.
GPU memory
Training runs close to the memory limit of the GPU, and the paper counts about 60 GB of an H100's 80 GB as available to a job. In early Llama 4 pre-training, NCCL's internal buffers took about 10 GB of that, spread over more than ten communicators (Section 7.2). The paper gives three reasons: NCCL allocates buffers and queue pairs up front for every peer, protocol and algorithm, whether or not they are used; each of its channels has its own set; and each channel's small metadata occupies a 2 MiB page of GPU memory, which adds up to more than 1 GB at 100K ranks.
NCCLX allocates and connects for an algorithm only when it is first used, creates channels only as they are needed, and packs the metadata of many channels into shared pages with a slab allocator (one 2 MiB page holds the metadata of about 3,000 peers). Table 4 reports the effect on a 64K-GPU pre-training run:
| Feature | Setting | GPU memory, before and after |
|---|---|---|
| Lazy algorithm connect | NCCL_LAZY_CONNECT=1 | 8.09 GB to 6.65 GB |
| CTran lazy connect | NCCL_LAZY_CONNECT=1 | 6.65 GB to 5.86 GB |
| Slab allocator | NCCL_MEM_USE_SLAB_ALLOCATOR=1 | 5.09 GB to 4.7 GB |
| Lazy channel allocation | NCCL_LAZY_SETUP_CHANNEL=1 | 5.39 GB to 4.2 GB |
Only the first two rows chain. The other two start from different numbers, so the table cannot be added up into one total. The paper summarizes the result as a reduction of almost 2ร. It also notes that the same features keep the number of queue pairs under 2,000 per network card, and that lazy connection was shared with NVIDIA and has a counterpart in NCCL since release 2.22.
Finding the fault
One bad network card or GPU can hang or kill a whole job, and with several dimensions of parallelism a collective that stalls in one dimension makes collectives in the others wait. NCCL's own diagnostic subsystem does not separate the original failure from the ones it causes, and it does not record whether a collective's kernel has started (Section 7.3).
NCCLX traces every collective and every RDMA operation, and a Fault Analyzer works on the traces with two assumptions. The job has hung long enough that every collective that could finish has finished. And a collective whose kernel has not started on a rank is waiting, directly or not, for the collective that is running on that rank. From these it builds the dependencies between collectives and finds the one that waits on nothing else. That is the first failure.
The paper gives two cases. In an 8K-GPU job that kept restarting, the analyzer found that the last AllReduce of the second data-parallel group was the first collective to fail, then found the host that first stopped sending, where one network card turned out to be faulty. In the other, a timeout that PyTorch reported only as a generic NCCL error was traced to one rank that never joined the last collective of its tensor-parallel group, because of a bug in the model code.
Testing without the GPUs
The scaling problems above appear only at sizes that are too expensive to test on real hardware. The authors built stand-ins for the CUDA and RDMA libraries and load them in place of the real ones, so that unmodified NCCL code runs on CPU servers (Section 7.5). With 32 processes per host, 3,072 CPU servers emulate 96K ranks. That is how they measured initialization at scale, and how they found the busy loops, the quadratic code and the overflowing TCP listen queues, which they fixed with reconnection retries and exponential back-off.
Results in one table
| What was measured | Result | Source |
|---|---|---|
| Latency of a steady Llama 4 training step | Up to 12% lower | Introduction |
| Creating the default process group, 96K GPUs | 265 s to 24 s, 11ร faster | Figure 21 |
| Point-to-point send, 1 MB to 128 MB | 1.09ร to 2.7ร faster than the copy-based send | Figure 10 |
| Tensor parallelism with overlap, one node | 1.57ร lower end-to-end latency | Figure 11 |
| Fault Tolerant AllReduce against NCCL AllReduce | Comparable with half the thread blocks, 9% to 18% lower at equal blocks | Figure 12 |
Decode time with AllToAllvDynamic | 16.65% to 83.3% lower than the CTran baseline | Table 3 |
| GPU memory used by NCCL | Almost 2ร lower | Table 4 |
Critical analysis
Strengths:
- It describes a system in production. NCCLX carried the training and serving of Llama 4, the code is public, and some of its changes were offered back to NCCL.
- One decision explains most of the gains. Running the collective on the CPU and letting the network card move the data removes the GPU threads, the intermediate buffers and the clear-to-send messages in one step. The custom collectives are then short programs on top.
- The layers are designed together. The allocator registers buffers, the scheduler places ranks by topology, the transport sets its limits by distance, and the model code changes its sharding and its buffering to match.
- The operational sections are rare. Startup time, memory held by the library, fault localization and testing at scale are seldom written down in this detail.
Limitations:
- Most charts are normalized. The latency and bandwidth figures (7, 10, 11 and 12) show ratios. The paper gives no absolute bandwidth or latency, so the numbers cannot be compared with other systems.
- The headline training number has no experiment behind it. "Up to 12%" lower step latency appears in the introduction, with no table, figure or configuration.
- The text and the tables disagree in places. The text gives the decode-time gain as 15% to 80% where Table 3 runs from 16.65% to 83.3%. Two rows of that table's improvement column do not follow from their own times: 73.82 ms to 44.09 ms is 40.27% and is printed as 44.27, and 50.51 ms to 20.09 ms is 60.23% and is printed as 59.96. The text says lazy connection saves about 2.8 GB where the two rows of Table 4 add up to 2.23 GB. Table 2's caption describes 32 nodes of 8 ranks with messages under 128 KB, and the text around it describes 128 GPUs on 32 nodes with 8 MB messages.
- Distributed inference is compared with itself. The 83% is measured against a distributed baseline that gets slower with every added host. Against the single-node times in the same table, at a batch of 128, the picture is mixed: with one expert per token the distributed runs have a 37% to 44% lower decode time than one node, and with four experts per token they are level at 4 hosts and 18% to 22% higher at 8 and 16. Those percentages are this page's arithmetic from Table 3.
- Tensor-parallel overlap is measured on one node. The design exists to cross nodes, and the cross-node case is described but not benchmarked.
- Fault tolerance has no outcome number. The paper motivates hybrid sharding by the share of productive training time, and reports the latency of the AllReduce, not that share.
- NCCL is the only baseline. MSCCL, UCCL and DeepEP are discussed as related work and not measured.
- The settings are not given. The per-distance limits of DQPLB, which carry the flow control, appear only as "conservative" and "aggressive". The released code exposes them as configuration.
Where it fits
NCCLX stands between two lines of work. From HPC it takes the host-driven model of MPI, its collective algorithms and its one-sided windows. From the GPU libraries it keeps NCCL's interface and its place under PyTorch, and it aims at the latency that device-initiated libraries reach.
It also shows what the tensor parallelism of Megatron-LM costs at scale. Megatron-LM left a few collectives per layer on the critical path and kept each model-parallel group inside one server. Here the same collectives cross nodes, and the paper's answer is to overlap them with the multiply and keep them off the GPU's threads.
Related Reading
- Megatron-LM: the tensor parallelism whose AllGather and ReduceScatter this paper hides inside the matrix multiply
- Every ยตs Matters: shortens small AllReduce operations inside one node, the low-latency end of the same problem
- NanoFlow: overlaps compute, memory and network work inside one GPU when serving a model
- Switch Transformers: the mixture-of-experts routing that produces the AllToAll of the inference section
- HiveD: keeps whole nodes available to the tenants that reserved them, which topology-aware placement depends on
- Multi-GPU communication: NVLink, InfiniBand, RoCE and the collectives themselves
