Can Cross-Node Small-Message Communication Stay Resident?

About two thirds of the 13.6 us in an 8 B cross-node allreduce is control cost paid again on every operation. That cost can be paid once. An existing resident-proxy design already captures it, and two patches leave it only 1.6x from our upper bound.

During decode, a large model runs the whole network once per generated token. Tensor parallelism does one or two allreduces per layer. Expert parallelism does a dispatch and a combine, two all-to-alls per layer. Each message is a few KiB to a few tens of KiB. On a 100 Gb network the wire time for that size is a fraction of a microsecond, so almost all of the latency is fixed cost: kernel launch, proxy wakeup, buffer bounce, protocol handshake, and completion. Once the model is forced to span nodes, that fixed cost times dozens or hundreds of communications per step becomes a visible slice of decode latency. This post asks whether moving those costs from once per operation to once per epoch is a new system abstraction, and why the answer is no.

What one small cross-node message passes through

Start with the scale. Every measurement here is on NVIDIA EDR InfiniBand (100 Gb/s, ConnectX-5), with one GPU on each of two nodes. A CPU RDMA write ping-pong in host memory is about 1.15 us one way. That is the hardware floor of the wire, the switch, and the two NICs. Launching an empty kernel and synchronizing is about 6.4 us on one GPU. Replaying a single-node CUDA Graph and synchronizing is about 5.7 us. An 8 B device-to-host copy plus a sync is about 6 us. Every later number should be read against 1.15 us and these single-GPU costs.

NCCL is the de facto standard for GPU collectives (NVIDIA NCCL). Across nodes, the NCCL kernel on the GPU places data in a buffer and updates a counter. A host proxy thread polls that counter, and once it sees new work it posts an RDMA request to the NIC. The peer proxy sees the completion and notifies the peer kernel. Every communication walks this chain. CUDA Graphs capture a sequence of kernels and NCCL calls into one graph, so later steps replay the whole graph and the host launch cost becomes once per graph rather than once per operation (CUDA Graphs). Proxy wakeup, buffer handoff, and protocol synchronization are still paid every time.

GPUDirect RDMA (GDR) lets the NIC read and write GPU memory directly and skips the host bounce buffer (GPUDirect RDMA). The driver has to expose GPU pages to the NIC, either by loading nvidia_peermem into the RDMA peer-memory interface or by exporting DMA-BUF from the open kernel module and registering it. Without GDR, NCCL falls back to a host bounce buffer. The function is the same. The extra cost is a PCIe round trip.

Two lines of work already do this

Two lines are already moving fixed cost off the critical path.

The first is GPU-initiated networking. NVSHMEM organizes GPU memory into a symmetric partitioned global address space, and a kernel can call put and signal directly. Its IBGDA transport lets a GPU thread build the NIC work request and ring the doorbell, so the CPU is off the path (NVSHMEM). DeepEP builds a low-latency MoE all-to-all on top of it, and the receiver only waits on a signal (DeepEP). Since 2.28, NCCL also has a device API. GIN (GPU-Initiated Networking) lets a kernel post a network put and a signal directly. A symmetric window is registered once, and later operations only write the signal (NCCL Device API). On H100 with 400 Gb CX-7, published 8 B one-way numbers are 8.35 to 9.0 us for GIN and 8.0 to 12.15 us for NVSHMEM (Hamidouche et al., 2025). Those figures are half of a ping-pong round trip, and the paper does not say whether each round relaunches the kernel. The shared price of this line is GDR. A GPU can drive the NIC only if the NIC can reach GPU memory.

The second line is a resident host proxy. MSCCL++ splits communication into persistent channels such as PortChannel and MemoryChannel. The GPU writes a trigger into a host-visible FIFO. A spinning CPU proxy (ProxyService) reads it and posts an RDMA write and an atomic signal against pre-registered memory. The GPU spins on the semaphore (MSCCL++). The UCCL line also puts the transport on a resident CPU core, as a substitute for IBGDA (UCCL). What they remove is proxy wakeup and descriptor construction. What they keep is one GPU-to-host PCIe write and the NIC round trip. On the HPC side, the persistent collectives and partitioned communication in MPI 4.0 already express the same semantics: install the request once, then send one ready token per partition (MPI Forum, 2021).

Before any new code, the defining ability (install once, then only consume tokens) was already expressed by several existing primitives. The remaining question is empirical. On an ordinary cluster, do these primitives run, and how far do they get?

TL;DR

Cross-node collectives at decode size are dominated by a fixed cost paid on every call. The natural guess is to install the communication once, and after that only consume a ready token and a completion token.

On an EDR InfiniBand cluster without GPUDirect RDMA, a CUDA Graph NCCL allreduce of 8 B takes 13.6 us. The hardware floor is 1.15 us. About 8.9 us is NCCL’s per-operation control. A resident host proxy can reach 4.7 us, so the gain is real. MSCCL++ PortChannel has the same shape, but it will not start, because it treats pinned host memory as GPU memory. Two patches later it is 9.8 us on V100, 1.6 times our bound. The rest is fences, buffer classification, and polling. One NCCL switch, shared buffers off, also halves 8 B send/recv.

An end-to-end estimate from the measured curves saves more than 10% only for small-batch cross-node MoE expert-parallel decode. The resident proxy is already there. What remains is tuning.

The conjecture we can test

The starting point was a number from an earlier measurement. On one pair of A40 nodes, a CUDA Graph NCCL send/recv was 37.7 us one way, while an RDMA write in host memory was 1.13 us. The wire was 3% of the time. The rest sat in the endpoints. If a communication that repeats inside an epoch, with a fixed shape and a fixed peer, becomes a resident object, then queue-pair setup, memory registration, the proxy thread, and a persistent kernel are paid once at install time. After that, each use is one ready signal and one completion signal, and most of the fixed cost on the decode path can be amortized. The optimistic bound at the time was 10% to 18% of a decode step in the most favorable case.

The judgment then was that 10% to 18% is worth measuring no matter who it belongs to. Missing GDR might make part of the cost a per-message bounce that residency cannot amortize. Another part is protocol and proxy, which is what GIN and IBGDA aim at. The net gain of this mechanism over the strongest existing combination is the gap between a resident state machine and those existing paths on the same hardware. That gap had no data.

What we measured

There are two pairs of machines. One pair is two servers, each with one NVIDIA A40 (Ampere, 48 GB), EDR InfiniBand, NCCL 2.28.8, and CUDA 13.0. The other pair is two servers, each with one V100 (Volta, 32 GB), on the same EDR fabric. CUDA 13 no longer supports Volta, so that pair uses CUDA 12.4, NCCL 2.20.5, and UCX 1.16. On both pairs NCCL reports GDR disabled in every configuration, and forcing the GDR level does not change that. A survey of all 20 GPU nodes in the cluster found ConnectX-5 NICs everywhere. nvidia_peermem is installed on every node and loaded on none. The datacenter GPUs use the closed kernel module, so they have no DMA-BUF. The one card that does have DMA-BUF has no second node to pair with. Every conclusion here is limited to an environment without GDR.

We wrote a resident host-proxy upper bound, called B4 below. A persistent kernel writes a flag mapped into host memory. A spinning CPU thread sees it and posts one RDMA write to already-registered memory. The peer’s persistent kernel polls the arriving flag and data. This is the simplest form of the MSCCL++ PortChannel family. In a second round we added byte-for-byte checks on the receiver and acquire-semantics polling. Every iteration was zero errors.

The metric is the median one-way time, a ping-pong round trip divided by two, for messages from 8 B to 64 KiB. The GPU side uses globaltimer or CUDA events. Each point is 1000 iterations, with one warmup and then five recorded runs. The controls are NCCL allreduce and send/recv under Graph and eager, plus NCCL knobs for protocol, channel count, shared buffers, launch mode, proxy batching, and the network plugin. We also ran UCX put, tag, and active message on GPU memory, MSCCL++ PortChannel with ProxyService, NCCL GIN, NVSHMEM with ibrc, ibgda, and ucx remote transports, and the UCCL p2p engine. The measured latency curves were then dropped into several decode shapes for a closed-form end-to-end estimate.

Where the 13.6 us goes

Per-operation NCCL path versus a resident proxy, and a breakdown of the 13.6 us in an 8 B graph allreduce.
Figure 1. Per-operation NCCL control is about two thirds of an 8 B allreduce. A resident proxy keeps one PCIe write and the NIC round trip.

The left side of Figure 1 is the two paths. The top row is NCCL after a CUDA Graph capture. Each operation still enters the kernel, waits for the proxy to notice new work, bounces through host memory when there is no GDR, posts the network request, synchronizes the protocol, and completes in the peer kernel. The bottom row is the resident proxy. The kernel and the proxy thread are already resident. Each operation is one flag write, one RDMA write, and one poll.

The right side breaks down the 13.6 us of an 8 B Graph allreduce on the A40 pair. The gray 1.15 us is the hardware floor. The blue 2.7 us is the PCIe round trip between GPU and host without GDR, from B4’s flag-only 3.76 us minus the floor. The green 0.9 us is the extra cost of moving a 16 B payload through mapped memory. The remaining hatched 8.85 us is NCCL’s own per-operation control: proxy wakeup, bounce, protocol sync, and kernel entry. B4’s 4.66 us sits at the dashed line. This supports the first-round judgment. About 66% of small-message latency is control cost that can move from once per operation to once per epoch, and a CUDA Graph does not get it, because a Graph only amortizes the host launch. Eager send/recv is more extreme. Of 27.6 us, launch plus sync is only 3.3 us. The other 19.6 us sits on kernel boundaries, the proxy, and the bounce.

Two surprises showed up here. The earlier 37.7 us depends on the node pair. The same image and the same method on another pair is 24.3 us, so a number carried across pieces of work has to be remeasured on the pair it will be compared against. Most NCCL knobs do nothing at 8 B. The LL protocol matches the default. LL128 and Simple are slower. Grouped launch, proxy batching, and a different network plugin do not move the number. At 64 KiB, B4 is 18.0 us and the NCCL Graph allreduce is 37.0 us, so the gap shrinks to about 2x. Fixed-cost dominance holds only below a few KiB.

At this point the judgment passed the first gate. The gain is real, and it is large. The second round asks whether an existing mechanism already takes that gain.

How much existing mechanisms already take

One-way latency ladder on a V100 pair without GPUDirect RDMA. Patched MSCCL++ PortChannel is 1.59x the resident-proxy bound.
Figure 2. Without GDR, the isomorphic MSCCL++ PortChannel is 1.6x from the resident-proxy bound after two patches. Every stack that requires GDR fails to initialize.

Figure 2 is every mechanism that ran on the V100 pair. The axis is the median one-way time at 8 to 16 B. The parentheses are the ratio to the validated B4. The corner lists implementations that did not initialize. Gray is the 1.1 us hardware floor. Green is two fence variants of B4. Blue is MSCCL++. Red is UCX and NCCL.

The GDR-dependent family cannot initialize in this environment. NCCL GIN reports neither peermem nor DMA-BUF, and both the device API and global GIN support are 0. NVSHMEM ibrc fails to build its transport table after the DMA-BUF check fails. ibgda reports that the peer GPU is not reachable. ucx reports Bad address while registering the GPU heap. The UCCL p2p engine only offers the two GPU-memory registration paths that need peermem or DMA-BUF, and it has no host-bounce mode. DeepEP sits on NVSHMEM and is unavailable with it. These failures point at GPU-memory registration, not at the communication mechanism, so we record them as unmeasurable here rather than as open space.

MSCCL++, which has the same shape as B4, also failed at first, with an error that nvidia_peermem was not loaded. The source shows that the design does not require GDR. It classifies a buffer by asking which device owns the address, and pinned host memory from cudaHostAlloc is classified as GPU memory, so registration takes the peermem branch. Its semaphore also lives in GPU memory, and the proxy’s RDMA atomic on that semaphore needs GDR too. Two changes fix it. Register host memory as host memory according to the CUDA pointer attributes, and put the semaphore in mapped pinned host memory. PortChannel plus ProxyService then runs, with zero errors on a byte-for-byte check, 9.8 us one way at 16 B, and 31.3 us at 64 KiB.

The important distance in Figure 2 is that row against B4: 6.1 versus 9.8 us, a factor of 1.59. B4 itself, rewritten with a per-thread system fence, is 7.8 us, and against that conservative version MSCCL++ is only 1.2x slower. Of the 3.6 us between them, a substantial part is the fence. The rest is buffer classification, proxy polling, and how the signal is posted, which are implementation differences inside one design. On the NCCL side, NCCL_NET_SHARED_BUFFERS=0 gives each peer its own bounce buffer. An 8 B send/recv then drops from the default 23.5 us to 11.6 us. One switch halves it, and it is the largest NCCL knob we found. No knob pushes allreduce below about 15 us. A UCX active message on GPU memory is 10.6 us.

How much of the saved time reaches decode

Analytic decode-step savings from replacing NCCL with the resident proxy. Only small-batch cross-node MoE expert parallelism clears 10 percent.
Figure 3. A closed-form estimate on the measured curves. Only small-batch cross-node MoE expert-parallel decode saves more than 10% of the step.

Figure 3 substitutes the measured A40 curves for Graph NCCL and for B4 into several decode shapes: Llama-3.1-8B and 70B with tensor parallelism across two nodes, and Qwen3-30B-A3B with 8-way expert parallelism across two nodes, at batch 1, 8, and 64. The model shape sets the number of communications and the message size per step. Compute time comes from a calibration on the same class of A40. The vertical axis is the fraction of the step saved by replacing NCCL with B4, and B4’s advantage is capped at 64 KiB, because we did not measure it still leading on larger messages. This is a closed-form estimate, not an end-to-end run. Expert-parallel all-to-all is approximated as one point-to-point times the number of peers, which biases the savings upward.

Only two kinds of cell cross the 10% line: Qwen3-30B-A3B expert parallelism saves 22.9% at batch 1 and 11.3% at batch 8, and Llama-3.1-8B tensor parallelism of 4 across nodes saves 12.4% at batch 8. The 70B models are compute-heavy, and the saved fraction stays under 4%. At batch 64 the message exceeds 64 KiB and the gain falls below 5%. Pipeline parallelism across nodes has only two communications per step, so the fraction is 0 and is not drawn. The baseline here is the B4 upper bound. Against MSCCL++, which already exists, or against NCCL with the knob turned, the per-operation room shrinks from about 9 us to about 3.6 us. A normal scheduler also tries not to put tensor parallelism across nodes. A search on the same class of A40 chooses tensor parallelism inside a node and data parallelism across nodes, so the tensor-parallel columns themselves sit away from a real deployment.

Why the conjecture does not hold

The first half holds. On an ordinary InfiniBand cluster without GDR, about two thirds of an 8 B cross-node communication is control paid on every operation, and that control can be paid once. A CUDA Graph does not get it. A resident proxy does. The second half does not hold. This is not a missing capability. MSCCL++ PortChannel already has the resident proxy, the trigger FIFO, and the atomic signal. It failed to run here because of the environment and one buffer-classification detail. After that detail is patched, it sits 1.2 to 1.6 times from our bound. The rest is fences, polling, and adaptation. That is a tuning pass, not a new abstraction. The end-to-end saving is also narrow. It shows up in small-batch cross-node expert parallelism.

A few measurements are missing. With GDR we never compared GIN or IBGDA with B4 on this hardware. The published numbers come from a faster network and a faster GPU, and the timing boundary is unclear, so we cannot say whether the device-initiated path has reached this bound. Those two lines were built for this problem, so the unmeasured part is more likely to shrink the remaining room than to enlarge it. The patched MSCCL++ payload lives in host memory, so every GPU read and write crosses PCIe. That makes the number slower, which favors our conclusion, and 1.6 times is an upper bound on the gap. The end-to-end estimate treats all-to-all as point-to-point and biases the savings upward. The V100 and A40 numbers are not on the same pair of nodes. Comparisons stay inside one pair. The first-round B4 had no byte-for-byte check. After the second round added one, the number moved from 6.19 to 6.14 us. The earlier 4.7 us is not an artifact. It is an A40 number, and it cannot be set next to MSCCL++ on V100.

What these measurements support directly

The bottleneck of small cross-node messages is endpoint software, and mostly the control steps paid again on every operation. On EDR InfiniBand without GDR, the hardware is 1.15 us of an 8 B Graph allreduce of 13.6 us. The PCIe round trip plus payload is 3.6 us. The remaining 8.9 us is proxy wakeup, bounce, and protocol sync. A CUDA Graph only removes the host launch, and it never touches this part. The claim holds below a few KiB. At 64 KiB the gap is already down to about 2x.

On a cluster without GDR, the newer device-initiated stacks fail as a group, and the error does not say that the cause is driver configuration. Ordinary NCCL silently falls back to a host bounce, so day-to-day training and inference look fine. GIN, NVSHMEM, UCCL, and DeepEP fail during initialization. Before arguing about the mechanism, check whether nvidia_peermem is loaded, whether the driver is the open kernel module, and whether the CUDA properties report GDR and DMA-BUF support.

An existing system that matches your design, and that only fails to run because of the environment, is still your baseline. MSCCL++ looked like an unsupported environment. It was a buffer-classification detail, and two edits brought it next to our bound. Stopping at the error would have turned a tuning problem into a missing capability. The other case is different. When the existing mechanism itself depends on the missing capability, as IBGDA does on GDR, its absence only means neither side was measured. It does not assign the space to a new design.

For small send/recv, try NCCL_NET_SHARED_BUFFERS=0 first. On our V100 pair it cuts 8 B one-way time from 23.5 us to 11.6 us, more than every other knob combined. Allreduce is insensitive to it. The spread across node pairs (37.7 versus 24.3 us under the same method) is also larger than most knobs, so a comparison of designs has to hold the node pair fixed.

The end-to-end value of residency depends on the ratio of communications per step to compute, not on how many times a single latency improves. Saving 2x to 3x per operation is 2% to 4% on a 70B tensor-parallel decode, and 11% to 23% on small-batch cross-node expert parallelism. Count the communications and the message sizes of the target parallel scheme first, then multiply by the measured per-operation gap.

The mechanism is already there

Only small-batch cross-node MoE expert-parallel decode saves more than 10%. The resident proxy is already there. What remains is tuning.