A distributed inference service moves two very different kinds of data between GPUs. Inside a model instance, tensor-parallel GPUs exchange partial results on every layer. Every GPU takes part at the same moment, in the same order, many times per token. Between instances, a prefill server hands a request's KV cache to a decode server, or a server loads a cached prefix from storage. That traffic is point to point, irregular in size and timing, and its peers come and go as the service scales.
NCCL was built for the first kind. NIXL, the NVIDIA Inference Xfer Library from the Dynamo project, was built for the second. This article explains the programming model of each, walks through NIXL's API lifecycle with real code, sizes a KV hand-off, and covers failure handling. For NCCL internals, read NCCL in depth. For the physics of KV transfers, read peer-to-peer KV transfer. For a full serving setup, read disaggregated serving.
NCCL's model: a fixed group in lock-step
An NCCL program starts by creating a communicator. One process calls ncclGetUniqueId, the ID is shared out of band, and every rank calls ncclCommInitRank(&comm, nranks, id, rank). From then on the membership is fixed. Collectives such as ncclAllReduce are enqueued on a CUDA stream, and every rank must make the same call, in the same order, with matching sizes. If one rank skips a call, or dies, the others wait forever inside a kernel.
NCCL also has point-to-point calls. ncclSend and ncclRecv are two-sided: both ends must post matching calls, normally inside ncclGroupStart / ncclGroupEnd so that sends and receives do not deadlock. That is fine for pipeline parallelism, where the schedule is known in advance. It is a poor fit for KV hand-offs. The decode server does not know in advance which prefill server it will hear from, adding a new server means building a new communicator, and a crashed peer can take the whole communicator with it. Recovery means ncclCommAbort and rebuilding.
NIXL's model: agents, memory and metadata
NIXL replaces the communicator with an agent, one per process, identified by a globally unique name such as prefill-0. An agent manages three things.
- Memory sections: address ranges registered with the agent. The segment types are DRAM, VRAM, NVMe-oF, object storage and files, so a GPU buffer and an S3 object are described the same way.
- Backends: transfer engines loaded as plug-ins. UCX is the default and covers NVLink, InfiniBand, RoCE and TCP. As of October 2026 the repository also lists GDS and GDS_MT, POSIX, OBJ, AZURE_BLOB, HF3FS, MOONCAKE, GUSLI, UCCL, GPUNETIO and LIBFABRIC. When a buffer is registered with several backends, the agent chooses one for each transfer based on the memory types at both ends and the backends the remote agent has.
- Metadata: a serialized blob with each backend's connection information and the remote keys for registered memory. Agents exchange it over a side channel or through a central store such as etcd, then cache it. Loading metadata does not open a connection; that happens on first use or through
make_connection.
Transfers are one-sided from the application's point of view. The initiator names local descriptors, remote descriptors, the remote agent and an operation, READ or WRITE, and gets back a handle. The target posts nothing. It learns that a transfer finished through an optional notification message delivered after the data. Membership is dynamic: a new agent joins by publishing its metadata, and a failed one is removed by invalidating its metadata.
The NIXL lifecycle in code
Install it with pip install nixl, which brings the Python API and a bundled UCX. The lifecycle below follows the two-peer example in the NIXL repository. The target registers a tensor and publishes descriptors; the initiator reads them.
import torch
from nixl import nixl_agent, nixl_agent_config
# enable_prog_thread, enable_listen_thread, listen_port
config = nixl_agent_config(True, True, 5555)
agent = nixl_agent("decode-0", config) # name must be globally unique
kv = torch.empty((10, 16), dtype=torch.float16, device="cuda")
reg = agent.register_memory(kv) # do this once, at start-up
# 1. metadata: fetch the peer's, publish ours (side-channel mode)
agent.fetch_remote_metadata("prefill-0", peer_ip, 5555)
agent.send_local_metadata(peer_ip, 5555)
while not agent.check_remote_metadata("prefill-0"):
pass
# 2. descriptors: the peer sent its serialized descriptor list (e.g. via send_notif)
remote = agent.deserialize_descs(serialized_from_peer)
local = agent.get_xfer_descs([kv[i, :] for i in range(10)])
# 3. transfer: READ remote -> local, notify the peer when done
h = agent.initialize_xfer("READ", local, remote, "prefill-0", b"req-42:done")
state = agent.transfer(h)
while state != "DONE":
if state == "ERR":
raise RuntimeError("transfer failed via " + agent.query_xfer_backend(h))
state = agent.check_xfer_state(h) # "PROC" while in flight; non-blocking
# 4. teardown
agent.release_xfer_handle(h)
agent.remove_remote_agent("prefill-0")
agent.deregister_memory(reg)The prefill side waits with agent.check_remote_xfer_done("decode-0", b"req-42") or polls get_new_notifs(), and frees its KV blocks only after that. Two rules from the NIXL design notes are worth taping to the monitor. Register memory and exchange metadata during initialisation, never on the request path. A transfer handle can be reposted many times, but only one transfer per handle may be in flight at once.
Worked example: a 2,000-token KV hand-off
Size a real hand-off. A Llama-3-8B-shaped model has 32 layers and 8 KV heads of dimension 128. In FP16, one token's K and V across all layers take 2 × 32 × 8 × 128 × 2 bytes = 131,072 bytes, or 128 KiB. A 2,000-token prompt is therefore 250 MiB (262 MB). On one 400 Gb/s NIC (50 GB/s at line rate) that is about 5.2 ms. That is acceptable against a time-to-first-token budget measured in hundreds of milliseconds, provided the transfer actually runs at line rate.
The catch is the descriptor count. With paged attention in 16-token blocks, the prompt is 125 blocks. If the cache is laid out per layer, as in vLLM, each block is a separate region in each of 32 layers, giving 4,000 descriptors, or 8,000 if K and V are separate tensors. Each is only 64 KiB, or 32 KiB when K and V are split. These numbers assume TP = 1. With TP = 4, each GPU holds 2 of the 8 KV heads and sends a quarter of the data over its own NIC, so every descriptor shrinks to about 8 KiB per K or V. Validating 8,000 descriptors per request on the hot path costs real CPU time, and many tiny RDMA operations do not reach line rate.
NIXL's answer is to prepare once and select by index. At start-up, each side calls prep_xfer_dlist over every block of every layer and keeps the handle. Per request, the initiator calls make_prepped_xfer with just the indices of the blocks involved. NIXL can merge descriptors that are contiguous in memory, which is why allocating a request's blocks next to each other directly raises throughput.
# start-up: one entry per (layer, block), in a fixed order on both sides
local_side = agent.prep_xfer_dlist("NIXL_INIT_AGENT", all_local_blocks)
remote_side = agent.prep_xfer_dlist("prefill-0", all_remote_blocks)
def idx(layer, block):
return layer * NUM_BLOCKS + block
# per request: only indices travel
ids_local = [idx(l, b) for l in range(32) for b in decode_blocks]
ids_remote = [idx(l, b) for l in range(32) for b in prefill_blocks]
h = agent.make_prepped_xfer("READ", local_side, ids_local,
remote_side, ids_remote, b"req-42:done")One design choice remains: who initiates. In a pull design the decode server issues a READ once it has allocated blocks for the request. Decode never receives data it has no room for, and the prefill server stays passive, which keeps the prefill side simple. The cost is that the transfer cannot start until prefill has finished and decode has been chosen. In a push design the prefill server issues WRITE operations into blocks that decode reserved up front. That lets prefill send each layer as soon as it is computed, so the transfer overlaps the rest of the prompt's computation. The price is a reservation handshake before prefill starts, plus cleanup when a reserved decode slot is abandoned. Start with pull, since it is easier to make correct, and move to push only if measured transfer time is a significant share of time to first token.
Choosing, and running both
| Question | NCCL | NIXL |
|---|---|---|
| Traffic shape | Collectives and scheduled send/recv | Irregular one-sided reads and writes |
| Membership | Fixed per communicator | Agents join and leave through metadata |
| Who posts | Every participating rank | Only the initiator; the target can be notified |
| Memory types | GPU memory buffers | DRAM, VRAM, files, block, object storage |
| Failure blast radius | The whole communicator; abort and rebuild | One remote agent; invalidate it |
| Typical use | Tensor, data and pipeline parallelism | KV hand-off, KV offload and load, weight loading |
In practice you run both. Tensor-parallel all-reduces inside each instance stay on NCCL, which is unbeatable for lock-step collectives on NVLink. The KV path between instances, and between GPUs and a storage tier, goes through NIXL. Frameworks hide most of this. vLLM exposes NIXL through its NixlConnector KV connector, and Dynamo uses it for its KV routing and offload. You still need the model above to debug them. See Dynamo and UCX for the layers above and below.
Failure modes
- Silent TCP fallback. UCX picks the best transport it can find. A missing RDMA device, a container without the InfiniBand verbs libraries, or a restricted
UCX_TLSgives you a transfer that works but runs many times slower. Logquery_xfer_backendand measured bandwidth for the first transfer after start-up. - No GDRCopy. NIXL and UCX work without it, but the README notes it is needed for maximum performance. Small, latency-sensitive GPU transfers suffer most.
- Stale metadata after a restart. A crashed agent that restarts under the same name has new memory keys. Peers holding the old cache fail or read the wrong memory. Invalidate on failure, remove the remote agent, re-fetch, and give each incarnation a unique suffix.
- Freed while read. If prefill frees or reuses KV blocks before the decode side's notification arrives, decode reads another request's cache, with no error. Free only on the notification, and time out and reclaim blocks for decodes that never arrive.
- Reposting an active handle. Posting the same handle again before it reaches
DONEbreaks the one-active-transfer rule and can corrupt data. Use a handle pool, or create one handle per in-flight request. - NCCL and NIXL sharing NICs. KV transfers can collide with inter-node collectives on the same rails. Watch per-port counters during load tests and, if needed, put the two on different NICs or traffic classes.
- A hung collective. On the NCCL side, a dead rank still freezes the communicator. Run a watchdog that calls
ncclCommAbortand restarts the instance, rather than letting a decode server hang with its KV memory pinned.
Benchmarking and configuration
Benchmark the path before trusting it. The repository ships nixlbench, which coordinates its workers through etcd, for example ./nixlbench --etcd-endpoints http://localhost:2379 --backend UCX --initiator_seg_type VRAM. Run it between the exact node pairs that will serve traffic, with message sizes matching your per-block descriptor size (8 to 64 KiB in the example), not just large messages. Then compare three numbers: line rate, nixlbench at your block size, and the transfer time your serving framework logs. A gap between the second and third usually means descriptor preparation on the hot path or fragmented block allocation.
NIXL reads configuration from environment variables first, then from a TOML file named by NIXL_CONFIG_FILE, $HOME/.nixl.cfg or /etc/nixl.cfg. A syntax error in that file does not fall back to the next one, so validate it in CI.
What to do next
- List every GPU-to-GPU flow in your service and label it lock-step collective or point-to-point hand-off. Keep the first on NCCL and move the second to NIXL.
- Run the two-peer example from the NIXL repository between two real nodes and confirm, with
query_xfer_backend, that UCX chose RDMA rather than TCP. - Compute your KV bytes per token and descriptors per request with the formulas above, then benchmark
nixlbenchat that descriptor size. - Move all registration, metadata exchange and
prep_xfer_dlistcalls to start-up, and check that the request path only callsmake_prepped_xfer. - Make block freeing depend on the completion notification, with a timeout that reclaims blocks from decodes that never come.
- Kill a prefill agent during a load test and check that peers invalidate it, that no request reads stale memory, and that a restarted agent rejoins under a new name.
- Add NCCL watchdogs with
ncclCommAbortso that a dead tensor-parallel rank restarts its instance instead of hanging it.