UCX, Unified Communication X, is an open-source communication library that sits between applications such as MPI and the network hardware: InfiniBand, RoCE, TCP, shared memory and GPU interconnects. If you run CUDA-aware Open MPI, use NVIDIA's NIXL for moving KV cache between inference workers, or shuffle GPU dataframes with RAPIDS, UCX is very likely what is moving your bytes.
Most engineers meet UCX only through environment variables pasted from a forum post. This article explains the model underneath them: the layers, how a message is sent, how GPU memory travels, and how to find out which path your traffic actually took. The hardware side is covered in InfiniBand for GPU clusters and GPUDirect.
Why a communication framework exists
Every fast network exposes a different low-level interface. InfiniBand and RoCE use verbs, with queue pairs, memory registration and completion queues; shared memory needs entirely different code; GPU memory adds CUDA IPC and peer-to-peer rules. Writing an MPI library against each of these directly multiplies effort. UCX gives one API and chooses the best transport per peer at run time: shared memory or CUDA IPC for a process on the same node, RDMA for a remote one, TCP if nothing else works.
For ML infrastructure this matters in three places: CUDA-aware MPI in HPC-style training and simulation codes, collectives through UCC, and point-to-point transfers of large GPU buffers such as KV cache in disaggregated serving, where TensorRT-LLM's default KV-transfer path is NIXL with UCX as its backend. NCCL, by contrast, ships its own InfiniBand and socket transports and does not need UCX for collectives; see NCCL collectives.
Architecture: UCP, UCT and UCS
UCX has three main layers. UCT (transports) is a thin abstraction over each piece of hardware, offering three send flavours: short (data placed inline in the work request), bcopy (copied into a pre-registered bounce buffer) and zcopy (the hardware reads the user buffer directly, which requires registration). UCP (protocols) is what applications call: tag matching for MPI, streams, remote memory access (put, get), atomics and active messages. UCP picks transports, fragments large messages and spreads them across rails. UCS (services) provides configuration parsing, memory pools, data structures and logging, and the memory-hook component intercepts allocation and release calls so registration and memory-type caches stay correct when buffers are freed.
| Transport | Use | Notes |
|---|---|---|
rc_verbs, rc_mlx5 | reliable connected RDMA | best latency and bandwidth; one QP per peer |
dc_mlx5 | dynamically connected (NVIDIA NICs) | scales to many peers with less memory |
ud_verbs, ud_mlx5 | unreliable datagram | bootstrap and wire-up; UCX adds reliability |
tcp | sockets | always available; slowest |
posix, sysv, cma, xpmem | same-node shared memory | xpmem and knem need kernel modules |
cuda_copy | host-GPU copies | used for staging |
cuda_ipc | GPU-to-GPU on one node | uses NVLink or PCIe peer access |
gdr_copy | CPU access to GPU memory via GDRCopy | low-latency small transfers |
Contexts, workers, endpoints and progress
The UCP object model is short. A context holds global resources and the selected features. A worker is a progress engine with its own completion queues, normally one per thread. An endpoint is a connection from a worker to one remote worker, created from the remote worker's address, which you exchange out of band (MPI uses its runtime; you might use a TCP socket or a key-value store).
Operations are non-blocking and return a request. Nothing guarantees progress unless someone calls ucp_worker_progress. This is the single most common bug in hand-written UCX code: a send that never completes because no thread is progressing the worker. A minimal tag send of a GPU buffer looks like this:
ucp_params_t params = { .field_mask = UCP_PARAM_FIELD_FEATURES,
.features = UCP_FEATURE_TAG };
ucp_config_t *cfg; ucp_context_h ctx; ucp_worker_h wrk; ucp_ep_h ep;
ucp_config_read(NULL, NULL, &cfg); /* reads UCX_* env vars */
ucp_init(¶ms, cfg, &ctx);
ucp_config_release(cfg);
ucp_worker_params_t wp = { .field_mask = UCP_WORKER_PARAM_FIELD_THREAD_MODE,
.thread_mode = UCS_THREAD_MODE_SINGLE };
ucp_worker_create(ctx, &wp, &wrk);
/* ucp_worker_get_address(wrk, &addr, &len): send addr to the peer out of band */
ucp_ep_params_t ep_p = { .field_mask = UCP_EP_PARAM_FIELD_REMOTE_ADDRESS,
.address = peer_addr };
ucp_ep_create(wrk, &ep_p, &ep);
void *gpu_buf; cudaMalloc(&gpu_buf, nbytes);
ucp_request_param_t rp = { .op_attr_mask = UCP_OP_ATTR_FIELD_MEMORY_TYPE,
.memory_type = UCS_MEMORY_TYPE_CUDA };
ucs_status_ptr_t req = ucp_tag_send_nbx(ep, gpu_buf, nbytes, /*tag*/ 42, &rp);
if (UCS_PTR_IS_ERR(req)) {
fprintf(stderr, "send failed: %s\n", ucs_status_string(UCS_PTR_STATUS(req)));
} else if (req != NULL) { /* NULL means completed inline */
while (ucp_request_check_status(req) == UCS_INPROGRESS)
ucp_worker_progress(wrk); /* nothing moves without this */
ucp_request_free(req);
}Declaring the memory type is optional, since UCX can detect it, but it skips a lookup on every call. The receiver posts ucp_tag_recv_nbx with a tag and mask and progresses its own worker the same way. Production code uses completion callbacks (UCP_OP_ATTR_FIELD_CALLBACK) and a progress thread or event loop rather than spinning.
Threading follows from the same model. A worker has one of three thread modes: UCS_THREAD_MODE_SINGLE (only one thread ever uses it), UCS_THREAD_MODE_SERIALIZED (several threads, but the application guarantees one at a time) and UCS_THREAD_MODE_MULTI (concurrent access with internal locking you pay for on every call). The usual high-performance design gives each communication thread its own single-mode worker and its own endpoints, so no locks are needed. When the application cannot spare a core for busy-polling, request UCP_FEATURE_WAKEUP and sleep on the worker's event file descriptor instead. Endpoint creation is not free either: wire-up exchanges transport addresses, so create endpoints once at start-up and reuse them rather than connecting per message.
Eager, rendezvous and multi-rail
UCP chooses a protocol by message size. Small messages go eager: data travels with the header and the receiver copies it into the posted buffer, or into an unexpected-message queue if no receive is posted yet. Large messages use rendezvous: the sender transmits a ready-to-send control message containing a memory key, and once the receiver matches it, data moves by RDMA, typically a get issued by the receiver straight from the sender's buffer, followed by an acknowledgement. Rendezvous avoids copies and unbounded buffering, at the cost of an extra round trip.
The switch point is UCX_RNDV_THRESH. By default UCX computes it from the transports present; you can override it, but measure first. UCX_RNDV_SCHEME selects get, put or automatic choice. Above the rendezvous threshold, large messages are split across up to UCX_MAX_RNDV_RAILS devices; the default uses the two best devices and the documented maximum is four.
How GPU memory travels
UCP accepts GPU pointers anywhere it accepts host pointers for tagged, stream and active-message operations (RMA and atomics have partial support). What happens next depends on where the peer is and what the system supports:
- Same node, other GPU:
cuda_ipcmaps the peer's allocation and the GPU copies directly over NVLink or PCIe. No host memory is involved. - Other node, GPUDirect RDMA available: the NIC reads and writes GPU memory directly. This needs the peer-memory kernel module or dmabuf support in a recent kernel and driver, and a sensible PCIe topology.
- Other node, no GPUDirect RDMA: UCX stages data through pinned host buffers with
cuda_copy, pipelining copies and network sends. It works, but bandwidth and latency suffer, and it consumes CPU memory bandwidth. - Small messages from host code:
gdr_copylets the CPU read or write GPU memory through a BAR mapping, avoiding a CUDA copy launch for tiny payloads.
Topology matters as much as software. A GPU should talk through the NIC on its own PCIe switch or NUMA node; crossing the CPU interconnect can sharply reduce effective bandwidth. Pin the pairing with UCX_NET_DEVICES (for example mlx5_0:1) per local rank, using nvidia-smi topo -m to read the matrix. More on the NIC side in RoCE v2.
Configuration and tools
Configuration is environment variables; ucx_info -c prints every one with its current value. The ones worth knowing:
| Variable | What it does | When to touch it |
|---|---|---|
UCX_TLS | transports allowed; aliases sm, rc, ud, cuda, tcp; prefix ^ to deny | exclude a broken transport, force a path while debugging |
UCX_NET_DEVICES | restrict to specific NICs and ports | GPU-NIC affinity, multi-NIC nodes |
UCX_RNDV_THRESH | eager-to-rendezvous switch point | only after measuring |
UCX_MAX_RNDV_RAILS | devices used for one large message (default 2, max 4) | set 1 to force the NUMA-local NIC |
UCX_MEMTYPE_CACHE | cache of pointer-to-memory-type lookups | set n if GPU buffers are misdetected |
UCX_LOG_LEVEL | set to info to print selected transports | first step of any investigation |
UCX_PROTO_INFO | prints the protocol selection table; on versions where the newer protocol selection is not the default, also set UCX_PROTO_ENABLE=y | seeing thresholds and chosen protocols |
# what does this node have?
ucx_info -v # version and build flags (was it built with CUDA?)
ucx_info -d | grep -E "Transport|Device"
# two-node bandwidth test, GPU buffers on both sides
ucx_perftest # on node A (server)
ucx_perftest nodeA -t tag_bw -m cuda -s 67108864 -n 100 # on node B
# Open MPI: use UCX, local GPU-NIC pairing, print selection
mpirun -np 16 --mca pml ucx \
-x UCX_NET_DEVICES=mlx5_0:1 -x UCX_LOG_LEVEL=info \
-x UCX_PROTO_ENABLE=y -x UCX_PROTO_INFO=y ./app
Worked example: slow inter-node GPU transfers
A team moves a simulation-plus-training code to two 8-GPU nodes. Intra-node GPU transfers look fine, but inter-node 64 MB sends reach a small fraction of the NIC's line rate in ucx_perftest. The investigation, in order:
ucx_info -dlistscuda_copyandcuda_ipc, so UCX was built with CUDA support. Good.- Running with
UCX_PROTO_ENABLE=y UCX_PROTO_INFO=yshows the large-message protocol for CUDA memory is a staged pipeline through host memory, not a zero-copy get. GPUDirect RDMA is not in use. lsmodshows no peer-memory module loaded and the kernel lacks dmabuf support for this driver. The team loads the module and re-runs: the table now shows zero-copy rendezvous for CUDA memory.- Bandwidth improves but varies by rank.
nvidia-smi topo -mshows half the GPUs share a PCIe switch withmlx5_1, notmlx5_0. A launcher wrapper setsUCX_NET_DEVICESper local rank to the closest NIC, and per-rank results become uniform.
No thresholds were changed. Almost every UCX performance problem in GPU clusters is a path problem (wrong transport, staging, wrong NIC) rather than a tuning problem, so diagnose the path before touching numbers.
Failure modes
- No progress. Hangs in custom code are usually an un-progressed worker; add a progress thread or poll in the event loop.
- Silent TCP fallback. If RDMA devices are not visible in a container, UCX picks
tcpand the job runs slowly instead of failing. Log transport selection at start-up and alert on it. - Containers missing devices.
/dev/infinibandnot mounted, or a memlock limit too low for registration. Checkulimit -l. - Version skew. The UCX inside an MPI build differs from the system one;
ucx_info -vfrom inside the job tells you which loaded. - Misdetected memory types with custom allocators or memory pools; try
UCX_MEMTYPE_CACHE=nto confirm. - Copy-pasted tuning from another cluster that disables the transport you need.
Trade-offs
UCX buys portability and automatic path selection at the cost of opacity: the same call can take very different paths, and you only find out by asking. NCCL is the better choice for dense collectives on NVIDIA GPUs; UCX shines for point-to-point, irregular and RMA-style traffic, MPI codes and mixed CPU-GPU workloads. rc transports give the best per-peer performance but use memory per connection, while dc scales to large jobs. For the underlying RDMA concepts, see RDMA fundamentals.
What to do next
- Run
ucx_info -vanducx_info -don a compute node and inside your container; confirm CUDA transports and RDMA devices appear in both. - Run
ucx_perftestbetween two nodes with-m cudaand keep the result as your baseline. - Start jobs once with
UCX_LOG_LEVEL=infoandUCX_PROTO_INFO=y(plusUCX_PROTO_ENABLE=yif no table appears) and record which transports and protocols were chosen. - Verify GPUDirect RDMA (peer-memory module or dmabuf) if CUDA traffic crosses nodes.
- Map each GPU to its nearest NIC with
nvidia-smi topo -mand setUCX_NET_DEVICESper local rank. - Alert when a job falls back to
tcp. - Only then consider threshold or rail changes, one at a time, against the baseline.