HOCCL: Offloading Collective Communication from GPU Cores to Accelerate Distributed Training Yao Fei1 , Gongming Zhao1 , Hongli Xu1 , Jin Fang1 , Jiacheng Zhu1 Shuo Xu2 , Kun Huang3 , Zhuolong Yu1 1 University of Science and Technology of China, China 2 China Mobile (Suzhou) Software Technology Co., Ltd., China
arXiv:2609.34334v1 [cs.NI] 28 Sep 2026
Abstract
Model GPT-6.7B
Large language model training involves massive computation on GPU streaming multiprocessors (SMs), the primary compute units of GPUs. Since SMs host specialized accelerators such as Tensor Cores, their efficient utilization is critical to training efficiency. Unfortunately, existing collective communication systems compete with computation for SMs, as they consume SMs for communication-related data movement and synchronization operations. We observe that communication can, in principle, be driven by DMA engines, thereby eliminating SM involvement in communication. Based on this insight, we propose HOCCL, a zero-SM collective communication framework consisting of three components: a stream manager, a point-to-point (P2P) executor, and a collective scheduler. The stream manager preserves operator-level temporal ordering with other GPU kernels. The P2P executor enables zero-SM point-to-point communication, while the collective scheduler orchestrates P2P transfers to maximize bandwidth. Experiments show that HOCCL preserves near-peak communication performance, achieving within 3% of the state of the art on average, while eliminating communication occupancy on nearly 10% of total GPU SMs. By freeing SM resources for computation, HOCCL improves end-to-end training throughput by up to 5%.
1
3 Pengcheng Laboratory, China
DeepSeek-V2-Lite
SM Usage Peak Avg. Peak Avg.
TP 16 4.8 – –
PP 8 7.1 8 7.4
DP 8 4.1 8 2.8
EP – – 32 9.8
Table 1: Peak and time-averaged communication SM occupancy under different parallelization strategies during training on 32 H800 GPUs, each with 132 SMs. consumes a non-negligible number of SMs that would otherwise be available to computation. On a platform with 32 H800 GPUs (each H800 has 132 SMs), we use widely adopted parallelization strategies [42] to train two representative workloads: a GPT-style 6.7B dense model and an MoE model, DeepSeekV2-Lite [9]. Using PyTorch Profiler [37], we measure the runtime overhead of communication kernels and report their average/peak SM usage in Table 1. Although TP and EP exhibit the highest peak SM usage, with peaks of 16 and 32 SMs, respectively, they overlap little with computation and therefore do not contend with computation for SM resources. In contrast, PP and DP are explicitly designed to overlap with computation [13, 23, 47], so their SM occupancy correspondingly reduces the SM capacity available to compute kernels. During GPT-6.7B training, PP and DP together occupy an average of 11.2 (7.1 + 4.1) SMs over time, accounting for nearly 10% of the total GPU SMs on an H800. This level of SM occupancy can reduce the performance of compute kernels by more than 10% [4, 10, 22] and will translate into millions of dollars of additional training cost [8]. This naturally raises the question of whether computationoverlapped communication, such as PP and DP, can be implemented without SMs. Existing collective communication systems rely on SMs primarily because an RDMA-capable NIC (RNIC) cannot directly perform RDMA reads or writes on arbitrary user buffers [26]. Instead, data must first be placed in RDMA-registered memory accessible to the RNIC, and this data movement is typically carried out by SMs [32] (see details in Figure 1a). Therefore, a major challenge is how to eliminate the need for SMs to move data into RDMAregistered buffers. One natural way to bypass this issue, as exemplified by VCCL [7] (see details in Figure 1b), is to register user buffers with the RNIC at runtime, allowing the RNIC to access them
Introduction
In recent years, large language models (LLMs) have achieved remarkable success and gained widespread adoption, demonstrating expert-level performance across a variety of domains such as natural language understanding [1, 3] and code generation [20]. Training large language models involves massive computation [6, 12, 21], which is executed primarily on streaming multiprocessors (SMs), the primary execution units of GPUs [29]. Specifically, SMs integrate highly specialized hardware (e.g., Tensor Cores [25]) to provide sufficient computation capacity for deep learning. Therefore, fully utilizing SMs is critical to training efficiency. However, our measurements uncover an underappreciated inefficiency in traditional training systems: communication This work has been submitted to IEEE for possible publication. Copyright may be transferred without notice, after which this version may no longer be accessible.
1
directly without SMs. In practice, however, this approach faces important practical limitations. First, this approach is incompatible with advanced virtual memory management (VMM) mechanisms (e.g., PyTorch expandable segments [36]), which may map a contiguous virtual address range onto non-contiguous GPU physical memory. Buffers managed by such VMM mechanisms cannot be registered with the RNIC, forcing VCCL to forgo efficient memory management and potentially incurring more than 10% additional fragmentation overhead [14]. Second, in modern training frameworks such as Megatron-LM [42], frequently communicated tensors (e.g., activations, gradients) are dynamically allocated and freed at runtime [24] to save memory, so their addresses change unpredictably. This makes repeated RDMA buffer registration unavoidable. In practice, the overhead averages several milliseconds but can fluctuate significantly, sometimes reaching tens of milliseconds [43]. In our experiments, this overhead will lead to nearly 10% end-to-end training slowdown (see details in Section 6.3). We therefore turn to a more practical solution: offloading data movement from SMs to dedicated DMA (Direct Memory Access) engines. DMA engines are CPU-controlled hardware units for bulk GPU memory movement. In current training systems, DMA engines are typically used only for data loading at the beginning of training and result offloading at the end, leaving them lightly utilized during most of training. This suggests that repurposing DMA engines for communication is feasible from a resource-availability perspective. Doing so efficiently, however, is nontrivial and raises three key challenges. First, it introduces a CPU–GPU coordination challenge. DMA engines are CPU-driven, but the CPU does not directly observe the execution state needed for timely control decisions. In practice, the GPU remains the primary execution substrate, and the progress of both DMA operations and compute kernels is primarily reflected in GPU-side state. Efficiently coordinating the control plane between the CPU and GPU is therefore nontrivial. Second, it introduces a transfer latency challenge. Compared with SM-based copies, DMA engines incur much higher startup latency due to their hardware characteristics and the long control path from the CPU to the DMA engine. In our measurements, this startup latency is about 5× that of SM-based copying, which leads to more than 20% communication performance loss. Therefore, effectively hiding this latency is critical. Third, it introduces a communication scheduling challenge. Unlike SMs, DMA engines are fixed-function hardware and cannot process multiple transfer tasks in parallel. This forces multiple communication flows to share DMA service in a serialized manner, making collective communication scheduling particularly challenging: an improper DMA task order can lead to deadlock, while poor DMA task scheduling can cause bandwidth imbalance and long-tail latency. To overcome these challenges, we present HOCCL, a zeroSM collective communication system. HOCCL has three key
components. The stream manager builds a reliable and efficient synchronization mechanism through handshake-based coordination between the CPU and GPU, leveraging GPU memory-operation hardware. The point-to-point executor employs a multi-stage pipelined transfer mechanism to effectively amortize DMA startup overhead. The collective scheduler maximizes bandwidth utilization through topology-aware flow scheduling and flow-size-aware chunk scheduling, while supporting arbitrary user-defined collective communication patterns. The main contributions of this paper are as follows: • Through an analysis of execution traces from large-scale model training, we identify a previously underexplored factor that limits training efficiency: communication operations contend with computation for shared compute resources (i.e., SMs), thereby slowing down training and reducing overall throughput. • We design and implement HOCCL, an SM-free collective communication framework that combines three components: the stream manager provides efficient CPU– GPU synchronization, The point-to-point executor employs a multi-stage pipelined design to hide DMA startup latency, and the collective scheduler schedules flows and chunks to maximize bandwidth utilization. • We evaluate HOCCL through benchmarks and end-toend Megatron-LM training on a testbed with 32 NVIDIA H800 GPUs. Extensive experiments show that HOCCL preserves near-state-of-the-art communication performance, with less than a 3% loss compared to NCCL, while using no SMs for communication. In end-to-end training, this removes communication occupancy on nearly 10% of total GPU SMs and improves training throughput by up to 5%.
2 2.1
Background and Motivation Distributed Training Workloads
Training large models typically combines multiple parallelization strategies, such as data parallelism (DP) [18,41], pipeline parallelism (PP) [13, 23], and expert parallelism (EP) [11, 17]. These strategies improve scalability by partitioning parameters and activations across GPUs, but they also introduce substantial communication. We classify communication into two categories. We refer to the first as background traffic, such as PP and DP traffic, because it can often be overlapped with computation to hide its overhead. We refer to the second as foreground traffic, such as TP and EP traffic, because it lies closer to the computation and is harder to hide. Background Traffic in Distributed Training. Some communication traffic in distributed training can be viewed as background traffic, because it is typically designed to overlap with computation to hide communication overhead. PP 2
SM
① load data App Data
⑧notify CPU ③ notify ④ Create wqe ② store data ⑦ Writes cq ⑤ NIC read data
SM SM
RDMA Buffer
⑥
GPU Mem
(a) NCCL
NIC NIC sends data out
① Register GPU Memory CPU ② Create wqe ⑤ Writes cq
App Data
(RDMA Buffer)
③ NIC read data
① launch memcpy ⑧notify CPU ③notify DMA Engine Memory Controller ⑦ Writes cq ④ Create wqe ② copy data ⑤ NIC read data RDMA Bufer App Data
④
NIC NIC sends data out
NIC
⑥ NIC sends data out
GPU Mem
(b) VCCL
(c) HOCCL
Figure 1: Data-plane and control-plane differences among NCCL, VCCL, and HOCCL. Solid lines denote the data path, and dashed lines denote the control path. NCCL is a fully SM-centric communication mechanism. VCCL allows the NIC to directly access user data by registering user buffers as NIC-accessible memory regions. In contrast, HOCCL offloads communication tasks to DMA engines and other memory-operation hardware, thereby reducing the work performed by SMs.
GPT-6.7B
Forward compute time
DeepSeek-V2-Lite
Backward compute time
10 8 6 4 2 0 TP
PP
DP
EP
Layer compute time (ms)
Contention rate (%)
and DP are two representative examples. In PP, many schedules [19, 23, 39] overlap activation transfers and gradient communication for one micro-batch with the forward or backward computation of another. In DP, prior approaches such as FSDP [47] prefetch the parameters of the next layer via AllGather during the computation of the current layer, thereby hiding parameter communication behind computation. Foreground Traffic in Distributed Training. TP and EP communication occurs within a layer and lies on the critical execution path. Subsequent operators often cannot proceed until collectives such as AllReduce, token dispatch, or token combine complete. As a result, unlike PP or DP traffic, TP/EP communication offers little room for overlap and typically exhibits a more serialized computation–communication pattern. Accordingly, prior ML systems work has focused either on maximizing communication performance [15, 46] or on reducing execution time through fused computation– communication kernels [2, 35, 48]. In this work, we focus on background traffic. To quantify how much communication occupies SM resources while compute kernels are active and how this affects computation, we define the contention rate, which measures the fraction of SMs occupied by communication while compute kernels are active. Let t ∈ T denote time within a training step, where T is the time interval of that step. Let C(t) be an indicator function that equals 1 when compute kernels are active at time t, and 0 otherwise. For a traffic class x (e.g., TP, EP, PP, or DP), we define the contention rate Ix as R C(t) SMx (t) dt R Ix = T , (1) NSM T C(t) dt
12 11.4 10.6 9 7.7 7.1 0
5
10
15
Contention rate (%)
(a) Contention rate of each traffic class
(b) One transformer-layer compute time vs. contention rate
Figure 2: The left figure reports the contention rate Ix of each traffic class during compute. The right figure shows that both forward and backward transformer layer compute time increase noticeably as the contention rate rises. Since TP and EP communication does not overlap with computation, their contention rates satisfy ITP = IEP = 0. In contrast, PP and DP account for all contention in both models, causing Isum to exceed 10% in each case. Figure. 2b further shows the relationship between contention rate and the execution time of the main computation in training (i.e., a Transformer layer) based on our experiments. The results show that when Isum ≈ 10%, forward and backward computation slow down by 8.5% and 7.7%, respectively. It also indicates that, for background traffic, even when communication is largely hidden behind computation, it can still impose a non-negligible performance penalty by reducing the SM resources available to compute kernels.
where SMx (t) is the number of SMs occupied by traffic class x at time t, and NSM is theP total number of SMs on the GPU. We further define Isum = x Ix as the total contention rate across all traffic classes. Figure. 2a reports the contention rate Ix of each parallelization strategy x during training for two representative models.
2.2 2.2.1
Existing Communication Mechanisms SM-Based Communication Mechanism
Streaming Multiprocessors (SMs). SMs are the basic execution units of NVIDIA GPUs and determine their compu3
tational capacity. In practice, GPU programs (e.g., general matrix multiplication [34], Allgather [28]) are issued in the form of kernel functions, which are executed in parallel across these SMs. Each SM integrates several key hardware components, including warp schedulers for instruction dispatch, load/store units for memory operations, CUDA cores and Tensor Cores for computation and fast on-chip memory resources, most notably registers and shared memory. The availability of these remaining resources determines whether an incoming kernel can be scheduled on that SM. SMs’ Role in Mainstream GPU Communication Systems. We summarize NCCL’s communication workflow in Figure 1a. The communication kernels in NCCL are responsible both for moving data and for coordinating communication progress and synchronization with the CPU or RNICs. The dashed lines represent how the GPU, the CPU proxy, and the NIC exchange control information, while the solid lines indicate the direction of data movement: ① SMs reading data from the user buffer into SM-local registers or shared memory, and then ② writing the data into a staging RDMA buffer. Once the chunk has been staged, GPU threads ③ notify the CPU proxy that the corresponding buffer region is ready. The CPU proxy ④ then prepares NIC work queue entries (WQEs) and rings the doorbell to post the send request. The NIC subsequently ⑤ reads the staged data from the RDMA buffer and ⑥ transmits it over the network. After the transfer completes, the NIC ⑦ writes a completion entry to the completion queue (CQ). The CPU proxy detects the CQ update and notifies the GPU, allowing the communication kernel to continue.
frequently allocate and release communication tensors during execution [24]; their addresses therefore change across iterations, making registration/deregistration a recurring overhead rather than a one-time setup. This overhead is typically at millisecond scale and can spike to tens of milliseconds [43], which can lead to up to 10% end-to-end training slowdown (see details in Section 6.3).
2.2.2
3
2.3
DMA Engines
DMA Engines. A DMA engine (also called a copy engine in NVIDIA documentation) is a dedicated hardware component on the GPU for data movement. DMA engines are controlled by the CPU and can move data without consuming SM resources. They were originally introduced to overlap host-to-device memory transfers with GPU computation on the transferred data. In addition to CPU–GPU transfers, modern DMA engines also support high-bandwidth inter-GPU memory transfers over NVLink, with peak bandwidth reaching up to 900 GB/s [33]. DMA Usage in Recent NCCL. Recent NCCL releases have introduced copy-engine based collectives for intra-node NVLink-domain communication [27]. However, this mode is limited to intra-node transfers and does not support cross-node communication over RNIC. Since background communication in distributed training is in most cases cross-node traffic, NCCL’s intra-node CE mode does not apply to these scenarios. HOCCL is instead a collective communication system that primarily applies DMA engines to cross-node communication while also supporting intra-node communication.
Registered-Buffer-Based Communication Mechanism
3.1
VCCL [7] is a representative attempt to reduce SM involvement in communication. Instead of letting GPU communication kernels copy user data into a staging RDMA buffer, it exposes user buffers to the RNIC through runtime registration, so the RNIC can read payloads directly (Figure. 1b). VCCL follows a simpler workflow: ① the CPU issues a memory region (mr) registration request to the RNIC. After the registration completes, ② the CPU prepares NIC work queue entries (WQEs) and rings the doorbell to post the send request. The RNIC then ③ reads the data and ④ transmits it. Once the transfer completes, it ⑤ notifies the CPU of the completion. This mechanism improves the data path but leaves two practical constraints. The first is memory-layout compatibility: with advanced virtual memory management (VMM), such as PyTorch expandable segments [36], contiguous virtual addresses may correspond to non-contiguous physical pages, which are not directly suitable for RNIC registration. In that case, compatibility with VCCL requires systems to forgo more advanced memory management mechanisms and instead adopt less flexible alternatives, incurring over 10% memory overhead due to fragmentation [14]. The second is registration cost. Training stacks such as Megatron-LM [42]
System Design Overview of HOCCL
HOCCL addresses three requirements of zero-SM collective communication: (1) preserving CUDA stream ordering, (2) sustaining point-to-point transfer efficiency despite staging and DMA startup overheads, and (3) balancing NVLink, RDMA, DMA-engine utilization. As shown in Figure 3, HOCCL consists of three corresponding components: the stream manager, the point-to-point executor, and the collective scheduler. Stream Manager (Section 3.2) provides HOCCL’s controlplane foundation. The CPU thread acts as the primary control entity, whereas key control information originates on the GPU, including the completion of preceding kernels and the progress state of DMA operations. HOCCL bridges this gap with GPU memory-operation hardware, which enables the GPU to directly read from and write to CPU memory. The handshake synchronization protocol preserves the ordering between the CPU communication process and the preceding and subsequent dependent kernels on the GPU stream, while the flag pool and buffer state machine enable efficient coordination among the RNIC, the CPU proxy, and the DMA engine. 4
Stream Manager
Flag pool and buffer state machine
CPU-GPU handshake synchronization protocol
Collective Scheduler
Point-to-Point Executor Inter-node
Intra-node
Topo-aware Flow schedule
Symmetric GPU buffer
Runtime IPC-handle exchange
Flow-size aware chunk schedule
GPU SM
Stream
Chunk pipeline Multi-QP parallel
DMA engine
state machine idle
CPU Mem flag pool for staging buffer start_flag end_flag CPU-GPU handshake
CPU-GPU handshake
memory operation hardware pre-comm kernel
no-dependency kernel
0 2
RNIC processing
1
DMA copying
State machine for staging buffers
after-comm kernel
DMA
Figure 4: Operation of the stream manager. It leverages memory-operation hardware for low-latency CPU–GPU memory reads and writes, together with a flag pool in CPU memory, to exchange control information.
CPU-based cross-process synchronization
HardWare GPU memory-operation hardware
progress communication
CPU
GPU control exchange for DMA and RNIC coordination. To achieve this, the stream manager leverages low-level CUDA driver primitives (e.g., cuStreamWriteValue32 and cuStreamWaitValue32) to operate GPU memoryoperation hardware, which can exchange control messages with CPU memory without occupying SMs. Built on top of these primitives, the stream manager implements an explicit CPU–GPU handshake synchronization protocol.
RNIC
Figure 3: System overview of HOCCL. HOCCL consists of three components. The stream manager acts as the control hub and coordinates the CPU and GPU using GPU memoryoperation hardware. The point-to-point executor serves as the main execution engine and handles both intra-node and inter-node communication. The collective scheduler serves as the scheduling engine and employs a two-stage algorithm to realize topology-aware collective communication with high bandwidth utilization.
3.2
Stream Manager
CPU–GPU handshake synchronization protocol. As shown in Fig. 4, the CPU needs a start notification before initiating communication and must later propagate completion back to the GPU stream. The stream manager therefore maintains two notification flags in CPU memory for each GPU stream, allowing direct CPU access and GPU access through memoryoperation hardware: start_flag and finish_flag. Together, these flags form an explicit CPU–GPU handshake that ensures both sides observe every state transition correctly.
Point-to-Point Executor (Section 3.3) provides the zeroSM transfer primitive used by higher-level collectives. For RNIC-based inter-node communication, it selects an RNICaffine symmetric GPU buffer, uses chunk pipelining to amortize staging overhead, and adopts multi-QP parallelism to mitigate DMA startup latency. For intra-node communication, it uses a CPU-based runtime IPC-handle exchange mechanism and a CPU-based cross-process synchronization mechanism to enable direct DMA-based transfers.
Listing 1: Handshake synchronization mechanism // On the main thread // stream1 notifies the CPU to start cuStreamWriteValue32(stream1,start_flag ,START_VALUE); cuStreamWaitValue32(stream1,start_flag,0); // stream1 waits for the CPU until completion cuStreamWaitValue32(stream1,end_flag ,FINISH_VALUE); cuStreamWriteValue32(stream1,end_flag,0);
Collective Scheduler (Section 3.4) orchestrates point-topoint flows into efficient collective communication. First, the Topo-aware flow schedule decomposes a collective operation into topology-aware stages to hide the intra-node bandwidth overhead introduced by HOCCL’s additional DMA copies. Then, the flow-size-aware chunk schedule allocates DMAengine time slices across flows, allowing them to share the same DMA engine and avoid long-tail effects.
// On the CPU proxy thread while(start_flag != START_VALUE){} start_flag = 0; // Process DMA and RNIC tasks until completion end_flag = FINISH_VALUE; while(end_flag != 0){};
Traditional communication libraries launch communication as SM-resident kernels on GPU streams [30]. Once enqueued on a stream, these kernels inherit the stream’s synchronization semantics, which naturally enforce the correct ordering between communication and the compute kernels before and after it. HOCCL, by contrast, removes such kernels from SMs, making conventional kernel-to-kernel synchronization inapplicable. The stream manager must therefore preserve the ordering between communication and surrounding GPU kernels, while also enabling efficient CPU–
As shown in Listing 1, the protocol consists of a start path and a completion path, through which the CPU effectively takes over coordination of the GPU stream. On the start path, the GPU stream writes a request tag (START_VALUE) to 5
Recv-side GPU
Switch
DMA
User-data
Symmetric GPU Chunk Pipe
MultiQP
PCIe
NVSwitch
Send-side GPU
DMA NIC
NIC
Symmetric GPU
(a) Detailed design of inter-node communication. The staging buffer is placed on a symmetric GPU and divided across multiple QPs to support multi-stage pipelining. CPU Proxy 0 (Recv)
①-a Write memhandle
①-b Wait and read memhandle
CPU proxy 1
③-b Wait finish notify SHM ③-a Write finish notify (Send) -a Unblock ④-b Unblock ② Launch DMA ④GPU Stream GPU Stream Recv-side GPU
and wait finish notify
NVSwitch
Zero-SM Point-to-Point Executor
3.3.1
Inter-node Communication
Symmetric GPU buffer. Inter-node communication must traverse the RNIC, whereas DMA engines support cross-GPU data movement rather than RNIC-mediated transfer. HOCCL therefore allocates two RDMA staging buffers for each GPU on a symmetric peer GPU: one for outgoing data and one for incoming data. Here, a symmetric GPU is a GPU under the same PCIe switch and thus at the same distance from the RNIC attached to that switch, as shown in Fig. 5a. The DMA engine first moves data over NVSwitch into the staging buffer. Because the local GPU and its symmetric peer are equally distant from the RNIC, the subsequent RDMA transfer incurs no additional path overhead. DMA–RNIC chunk pipelining. The symmetric-buffer path adds a staging copy. To hide this cost, HOCCL overlaps GPU-side DMA transfers with NIC-side RDMA transfers. It organizes the staging region as a buffer pool, so different chunks can occupy different pipeline stages at the same time. On the sender side, the RNIC transmits ready staging buffers while the DMA engine fills free buffers with subsequent chunks. On the receiver side, the RNIC writes incoming data into receive buffers while the DMA engine drains completed buffers into the final user buffer. This decoupling amortizes DMA startup overhead and preserves high end-toend transfer efficiency. Figure 5a illustrates this design, where the arrows show the transmission path of each chunk. Multi-QP parallelism. Chunk pipelining overlaps DMA and RNIC work, but a single serial chunk stream would still expose CPU–GPU control latency and DMA-engine launch overhead. HOCCL therefore uses multiple RDMA queue pairs (QPs) per connection, as shown in Figure 5a. A QP encapsulates a pair of send and receive work queues together with completion signaling. By striping chunks across multiple QPs (e.g., two or four), HOCCL keeps multiple chunks in flight between each communicating pair. This deeper pipeline hides CPU–GPU control latency and DMA startup overhead behind in-flight chunks and improves RNIC bandwidth utilization.
Node 1
Node 0
3.3
Send-side DMA GPU
Node 0
(b) Detailed design of intra-node communication. HOCCL implements intra-node communication through direct DMA-engine-based transfers, while using the CPU for runtime address exchange and control signaling.
Figure 5: The upper figure illustrates the data path and pipelined design for inter-node communication, whereas the lower figure shows the CPU–GPU coordination process for intra-node communication. start_flag. The CPU proxy polls this flag, consumes the request once it becomes visible, and resets it to zero. The GPU stream waits for this reset as an explicit acknowledgement before proceeding. On the completion path, after the communication operation finishes, the CPU writes a completion tag (FINISH_VALUE) to end_flag to notify the GPU stream of completion. The GPU stream then resets this flag as the final acknowledgement, completing the end-to-end handshake. This workflow preserves correct stream-ordering semantics without SM-resident communication kernels, while making each state transition explicit and allowing notification state to be safely reused Flag pool and buffer-state coordination. Once the CPU takes over coordination of the GPU stream and initiates a communication operation, it must manage two asynchronous execution paths simultaneously: RNIC operations on the CPU side and DMA tasks on the GPU side. To coordinate these paths efficiently, HOCCL uses a GPU-CPU-shared flag pool together with a per-chunk buffer state machine. Each stagingbuffer chunk is associated with a flag whose value encodes its current state in a three-state machine: idle, DMA copying, and RNIC in progress as shown in Figure 4. These flags are shared between the CPU thread and the GPU stream, allowing both sides to track chunk ownership and execution progress through a unified control structure.
3.3.2
Intra-node Communication
Figure 5b shows HOCCL’s intra-node design. HOCCL enables direct send/recv between GPU buffers while resolving the runtime metadata dependencies required for correct communication. Before issuing a DMA write, the sender must obtain the receiver’s destination address; before consuming the data, the receiver must know that the transfer has completed. HOCCL handles both metadata exchange and completion notification through CPU-side inter-process messaging, avoiding both staging-buffer-based communication [32] and designs that require addresses to be known before communication begins [40]. As a result, the data path remains direct, 6
′ where fm is the sender-side NVLink forwarding copy on ′′ (u, ū), fm is the RDMA transfer, and fm is the receiver-side NVLink forwarding copy on (v̄, v). Let X X ′ ′′ Au = fm , Bv = fm , (3)
asynchronous, and free of SM involvement, enabled by a runtime IPC-handle exchange mechanism and a CPU-based cross-process synchronization mechanism. Runtime IPC-handle exchange. HOCCL uses POSIX semaphores [45] and shared memory to exchange the control metadata required by the direct path. For each unidirectional (rank, peer) pair, the receiver-side proxy exports the destination buffer as a cudaIpcMemHandle and publishes the handle in shared memory. The sender-side proxy opens the handle at runtime, resolves the receiver’s GPU address, and issues a DMA transfer directly into the receiver buffer. Because this exchange is performed by CPU proxy threads, a local send/recv call does not block the calling CPU thread while waiting for the peer to post the matching operation; instead, the proxy completes the exchange once the peer-side metadata becomes available. CPU-based cross-process synchronization. After the sender-side proxy observes that the DMA write has completed, it posts the completion semaphore associated with the receiver. The receiver-side proxy waits on this semaphore and marks the receive task as finished once the signal arrives. The stream manager can then notify the corresponding GPU stream that the receive operation has completed, allowing subsequent computation on that stream to proceed.
3.4
m∈Ou
m∈Iv
where Ou is the set of inter-node flows sent by gu , and Iv is the set of inter-node flows received by gv . For an entry T [u, v] = Fi , we have ( C[u, v] + Au + Bv , T [u, v] = C[u, v],
(u,v)∈DNVLink v=ū,
otherwise.
(4)
This notation follows the second matrix in Figure 6: each Fi is the effective traffic of the corresponding cell after adding forwarding copies. For example, cell 1 corresponds to the NVLink edge from g0 to g1 . It carries the native intra-node flow f1 , the senderside forwarding copies for g0 ’s inter-node flows f2 and f3 , and the receiver-side forwarding copies for flows whose final destination is g1 , namely f8 and f11 : ′′ F1 = T [0, 1] = f1 + f2′ + f3′ + f8′′ + f11 .
Collective Scheduler
(5)
Thus, although C contains only user-visible logical flows, T captures the actual NVLink and RDMA traffic that HOCCL must schedule. Topo-aware flow schedule. After constructing T , HOCCL schedules traffic over two types of communication domains: intra-node NVLink domains and inter-node RDMA domains. In Figure 6, for example, there are three domains: one NVLink domain containing g0 and g1 , one NVLink domain containing g2 and g3 , and one RDMA domain connecting the two nodes. Executing all entries of T at once can oversubscribe some endpoints while leaving other domains underutilized. This is especially problematic for HOCCL because an RDMA flow also induces NVLink forwarding traffic. Therefore, HOCCL decomposes T into the stage matrices shown on the right side of Figure 6: P = {T 1 , T 2 , . . . , T n }. (6)
Figure 6 illustrates the workflow of HOCCL’s collective scheduler. It first translates a user-level collective description into an origin traffic matrix C, whose entries are the logical flows f1 , . . . , f12 . It then expands C into the HOCCL traffic matrix T , whose entries F1 , . . . , F12 include the forwarding traffic introduced by HOCCL’s zero-SM inter-node path. Finally, it decomposes T into stage matrices T 1 , T 2 , . . . , T n and schedules the selected flows as chunks on each DMAengine work queue. Traffic matrix abstraction. The origin traffic matrix C ∈ FG×G represents the collective exactly as requested by the user, where G is the number of GPUs. Each non-empty entry C[u, v] is a logical flow from source GPU gu to destination GPU gv . In Figure 6, for example, C[0, 1] = f1 , C[0, 2] = f2 , and C[3, 2] = f12 . Orange boxes denote send/recv pairs within the NVLink domain, which communicate over the intra-node high-bandwidth NVLink fabric( i.e., the intra-node path described above). Blue boxes denote send/recv pairs within the RDMA domain, which communicate via the RNIC (i.e., the inter-node path described above). HOCCL must convert C into the actual traffic seen by the underlying communication substrates. Let DRDMA denote the set of inter-node entries and DNVLink denote the set of intranode NVLink entries. For the topology in Figure 6, each GPU gu has a symmetric GPU peer ū : 0̄ = 1, 1̄ = 0, 2̄ = 3, and 3̄ = 2. An inter-node logical flow fm = C[u, v] is realized as three pieces: ′ ′′ fm ⇒ fm + fm + fm , (2)
Each stage T s contains a subset of effective flows that can execute together. The decomposition must satisfy three requirements: completeness, per-stage degree bounds, and crossdomain balance. (1) Completeness. The staged matrices must cover the full HOCCL traffic matrix: n X
Ts = T
(element-wise).
(7)
s=1
(2) Per-stage degree constraint. To reduce RDMA congestion, HOCCL restricts each GPU to at most one outgoing 7
Topo-aware flow schedule
Dst g0 g1 g2 g3 Src
Collective communication description
g0
GroupStart() HoCCL_Send(dst1,ptr,size) HoCCL_Send(dst2,ptr,size) HoCCL_Recv(src1,ptr,size) HoCCL_Recv(src2,ptr,size) ... GroupEnd()
f1 f2 f3
g1 f4 g2 f7 f8
f5 f6 f9
g3 f10 f11 f12
F1 F2 F3
F1 1 F 1 2 F1 3
F5 F6 F7 F8 F9
F1 4
F10 F11 F12
F110 F111F112
F4
F1 7 F1 8
Origin traffic matrix C Traffic matrix abstraction
Flow-size-aware chunk schedule
GPU0
User send buffer RNIC
DMA Engine
Intra-node P2P buffer RDMA Recv buffer chunks
F1 5 F1 6
DMA engine work queue
T1
F1 9
F2 1 F2 2 F2 3 F2 5 F2 6
F2 4 F2 7 F2 8
...Tn
F2 9
F210 F211 F212
T2
GPU1
RDMA Send buffer chunks Intra-node P2P buffer
RNIC
User recv buffer
Figure 6: Collective scheduler of HOCCL. HOCCL converts a collective description into the origin traffic matrix C, expands C into the HOCCL traffic matrix T by adding forwarding traffic, decomposes T into topology-aware stage matrices T 1 , . . . , T n , and finally slices the selected flows into chunks for the DMA-engine work queue. Algorithm 1 Birkhoff-like decomposition for the HOCCL traffic matrix Require: HOCCL traffic matrix T with entries Fi Ensure: Per-stage matrices P = {T 1 , . . . , T n } 1: Initialize the residual matrix R ← T and P ← ∅. 2: Decompose native NVLink traffic in each NVLink domain into matching components. 3: while R still has unscheduled traffic do 4: Create an empty stage T s . 5: Extract a contention-free RDMA matching from the residual matrix. 6: Move the largest removable portion of this matching into T s and subtract it from R. 7: Add the corresponding sender-side and receiver-side forwarding copies to T s . 8: Fill the remaining NVLink slack with native NVLink matching components. 9: Append T s to P ; s ← s + 1. 10: end while
and one incoming RDMA flow in each stage: X I[T s [u, v] > 0] ≤ 1, ∀u, s, v: (u,v)∈DRDMA
X
I[T s [u, v] > 0] ≤ 1,
(8) ∀v, s.
u: (u,v)∈DRDMA
where I[·] is the indicator function. This constraint makes the blue RDMA entries in each stage contention-free and prevents a single GPU or NIC from becoming a short-term bottleneck. (3) Cross-domain balance. Since stages are executed serially, the duration of a stage is determined by the slowest communication domain. Let Bd be the bandwidth of domain d, and let Vd denote the set of matrix entries belonging to that domain. We define the workload of domain d in stage s as Lsd = max size(T s [u, v]) . (u,v)∈Vd
(9)
The normalized execution time of domain d is Lsd /Bd , so the stage duration is approximated by the slowest domain: τ s = max d
Lsd . Bd
It first initializes the residual matrix and pre-decomposes native NVLink traffic in each NVLink domain into local matching components (lines 1–2). It then repeatedly constructs a new stage by extracting a contention-free RDMA matching from the residual matrix and removing the largest feasible portion (lines 3–6). For each such stage, the algorithm adds the required forwarding copies and uses the remaining NVLink slack to pack local NVLink matching components when possible (lines 7–8). The completed stage is then appended to the output set (line 9). Overall, the algorithm follows the matching-peeling intuition of Birkhoff-von Neumann decomposition [5] while accounting for forwarding traffic and heterogeneous domain bandwidths.
(10)
The scheduler therefore seeks a decomposition that reduces the total stage time: min s n
{T }s=1
n X
τ s.
(11)
s=1
In Figure 6, the lower-right matrices T 1 , T 2 , . . . , T n illustrate this process: each stage contains a contention-free set of RDMA flows and enough NVLink traffic to use the remaining NVLink slack. Birkhoff-like decomposition algorithm. Algorithm 1 decomposes the HOCCL traffic matrix into a sequence of stages.
Flow-size-aware chunk scheduling. Stage-level schedul8
4
NCCL-CE
HOCCL
NCCL-1sm
NCCL-2sm
NCCL-2sm
NCCL-4sm
Intra-node NCCL-8sm
NCCL-8sm
NCCL-4sm
200 150 100 50 0 2 4 8 16 64 256 1024 Communication Size (MB)
(a) Intra-node P2P
Inter-node
50 40 30 20 10
2 4 8 16 64 256 1024 Communication Size (MB)
(b) Inter-node P2P
Figure 7: Point-to-point bandwidth. in existing training frameworks, we also build a PyTorchcompatible backend [38]. This backend is implemented as a lightweight Python extension with approximately 1K lines of C++ code, and its interface is designed to match PyTorch’s torch.distributed backend abstraction. As a result, HOCCL can be registered and invoked in PyTorch with minimal integration effort. In our current prototype, switching between NCCL and HOCCL requires only changing an environment variable.
Limitations
HOCCL has several limitations. First, HOCCL introduces additional NVLink traffic, because each RDMA transfer is accompanied by sender-side and receiver-side forwarding copies. In principle, this extra traffic may contend with latency-sensitive foreground communication such as TP, since both use the same scale-up fabric. In practice, however, the interference is typically limited. The reason is that HOCCL’s extra NVLink traffic is only serves to feed or drain RDMA transfers, and its sustainable rate is therefore fundamentally bounded by the scale-out bandwidth. Since modern GPU clusters usually provide substantially higher scale-up bandwidth than scale-out bandwidth (900 GB/s vs. 50GB/s), the additional NVLink traffic often consumes only a fraction of the available intra-node bandwidth. To further protect latencysensitive foreground traffic, HOCCL employs a switch mechanism that gives strict priority to foreground communication: whenever TP traffic is active, HOCCL temporarily suspends its background tasks and resumes them only after the foreground traffic completes. This design prevents HOCCL from extending the critical path of foreground communication. Second, HOCCL does not naturally support reduction-based operators such as AllReduce and ReduceScatter. These operators require communication to be tightly coupled with reduction computation, and achieving high performance typically relies on SM-side execution. As a result, such operators fall outside HOCCL’s target scope.
5
HOCCL NCCL-1sm
Algorithm Bandwidth (GB/s)
Algorithm Bandwidth (GB/s)
ing determines which flows are active in a stage, while chunklevel scheduling determines how they share the serialized DMA engine at the bottom of Figure 6. This layer is necessary because a single DMA engine may serve sender-side forwarding copies, receiver-side forwarding copies, and native intra-node transfers, but can execute only one DMA task at a time. A naive per-flow schedule would let large flows monopolize the engine, causing forwarding backlogs and long-tail stragglers. To avoid this, HOCCL divides each flow into fixedsize chunks and allocates DMA service in proportion to flow size. At the beginning of each stage, the scheduler estimates the byte volume of all flows mapped to the same DMA engine and constructs a weighted interleaving of chunks, so larger flows receive more service while smaller flows are still revisited regularly. As illustrated in Figure 6, a DMA engine may interleave chunks from native intra-node traffic, sender-side forwarding traffic, and receiver-side forwarding traffic within the same stage. This flow-size-aware policy keeps per-flow service rates aligned with demand, allowing flows in the same stage to complete at similar times.
6
Evaluation
Our evaluation tests three claims about HOCCL. First, offloading communication from SMs should preserve high point-to-point and collective bandwidth. Second, the additional CPU control work introduced by host-driven progress should remain off the training critical path. Third, eliminating communication-related SM occupancy should improve endto-end training throughput. We evaluate these claims with communication microbenchmarks, CPU-overhead measurements, and Megatron-LM training workloads.
6.1
Experiment Setup
Testbed. We run experiments on an H800 cluster. Each node contains eight NVIDIA H800 80 GB GPUs. GPUs within a node are connected by NVLink, with up to 400 GB/s unidirectional bandwidth. Each node also has eight NVIDIA Mellanox ConnectX-5 NICs, each with 50 GB/s link bandwidth; we use a one-to-one GPU–NIC mapping for inter-node communication. The software stack uses CUDA 12.2, NVIDIA driver 535.161.08, ConnectX-5 firmware 28.38.1902, and NCCL 2.29.7, which was the latest NCCL release available at the time of our experiments [32]. We use NVIDIA nccl-tests built against this NCCL version for communication microbenchmarks, and Megatron-LM with a PyTorchcompatible backend for end-to-end training. Baselines. We compare against NCCL and VCCL. NCCL is the production device-driven baseline, which launches communication kernels on GPU SMs. Since HOCCL’s objective is to remove this SM consumption, a comparison
Implementation
We implement HOCCL by extending the NCCL codebase [32]. The core HOCCL runtime consists of more than 3K lines of C++ code. To make HOCCL usable 9
6.2
Model Layers Hidden Heads Seq. Parallelism 6.7B 32 4096 32 2048 PP=4, DP=8 32B 60 6656 52 2048 TP=4, PP=4, DP=2
Communication Benchmark
Intra-node P2P. HOCCL achieves high intra-node bandwidth without communication SMs. As shown in Figure 7a, HOCCL increases from 19.6 GB/s at 2 MB to 192.2 GB/s at 1 GB. NCCL-CE exhibits similar intra-node performance with HOCCL. NCCL’s SM-driven bandwidth, in contrast, scales with the SMs allocated to its communication kernels: under 1/2/4/8 SMs, it plateaus at roughly 20/40/80/156 GB/s on medium and large messages. NCCL-8SM is faster on small transfers because SM kernels have lower startup latency (41.0 vs. 19.6 GB/s at 2 MB), but HOCCL catches up at 64– 128 MB and exceeds NCCL-8SM by 23% at 1 GB (192.2 vs. 156.3 GB/s). Inter-node P2P. HOCCL saturates the inter-node link for large transfers, but pays higher startup cost on small messages. Figure 7b shows that HOCCL scales from 13.8 GB/s at 2 MB to 48.8 GB/s at 1 GB, close to the effective 50 GB/s NIC bandwidth. NCCL again depends on the SM budget when constrained: NCCL-1SM and NCCL-2SM plateau around 16 and 31 GB/s, while 4 SMs are sufficient to nearly saturate the link and 8 SMs add little. HOCCL lags on small messages because DMA startup and chunk-pipeline setup are harder to amortize, but the gap disappears once the transfer becomes bandwidth dominated. Thus, for large inter-node P2P, HOCCL reaches NCCL’s best-case bandwidth while avoiding deviceside progress. AllGather. HOCCL preserves strong AllGather bandwidth when NCCL cannot dedicate many SMs to communication. Figures 8a and 8b report an 8-GPU intra-node AllGather and a 16-GPU cross-node AllGather. In the 8-GPU intra-node case, NCCL-CE is close to HOCCL across message sizes and reaches 188.7 GB/s at 1 GB compared with HOCCL’s 192.2 GB/s. HOCCL is also 2.4× faster than NCCL-4SM (80.2 GB/s) at this size, showing that DMA-based intra-node transfers can match NCCL-CE while avoiding communication SMs. The 16-GPU AllGather experiment is conducted on a separate testbed, where we use one GPU from each participating node. This configuration forces all communication onto the NIC/RDMA path rather than the intra-node NVLink path. We choose this setup because AllGather commonly appears in DP communication, and prior large-scale training systems often reserve the fast intra-node fabric for bandwidth-sensitive traffic such as TP while placing DP communication across nodes [16, 44, 47]. This setup therefore isolates the inter-node behavior most relevant to our target DP workload. At 1 GB, HOCCL reaches 48.2 GB/s, comparable to NCCL-4SM and NCCL-8SM at 47.9 and 48.3 GB/s, while NCCL-1SM and NCCL-2SM plateau near 16 and 30 GB/s. NCCL-CE is omitted because it is intra-node only. This result shows that HOCCL extends zero-SM communication to the cross-node AllGather pattern that dominates many DP workloads. All-to-All. HOCCL is particularly effective for large All-
Table 2: End-to-end training model configurations. only against default NCCL would hide the central resource trade-off. We therefore also evaluate NCCL under explicit communication-kernel SM budgets of 1, 2, 4, and 8 SMs, making the bandwidth–SM trade-off visible. For intra-node experiments covered by NCCL’s copy-engine mode, we also report NCCL-CE [27]; because this mode does not support cross-node RDMA transfers, it is not included in inter-node P2P or collective benchmarks that cross nodes. VCCL is the most relevant host-driven baseline because it allows the NIC to access user buffers directly after registration. We do not include VCCL in bandwidth-only microbenchmarks. Standard microbenchmarks repeatedly reuse the same communication buffers, which amortizes registration cost and can substantially overstate the benefit of user-buffer registration. Communication benchmarks. We use nccl-tests [31] to isolate communication performance. We evaluate point-to-point transfers and two collectives: AllGather and All-to-All. Intra-node P2P runs within a single node, while inter-node P2P uses one GPU pair across nodes. For each collective, we evaluate two configurations: Config 1 represents intra-node communication, and Config 2 represents collectives involving inter-node traffic. NCCL-CE is included for intra-node P2P and Config 1 collectives, where its copyengine mode applies; it is omitted for Config 2 because NCCLCE does not support cross-node transfers. We sweep message sizes from 2 MB to 1 GB and report algorithmic bandwidth. For NCCL, we additionally vary the number of SMs assigned to communication kernels. End-to-end workloads. We evaluate training with Megatron-LM on all 32 H800 GPUs using GPT-style decoderonly Transformer models. Table 2 summarizes the model and parallelism configurations. Both workloads use BF16 mixed precision and sequence length 2048. For throughput experiments, we use global batch sizes 32 and 128 for the 6.7B and 32B workloads, respectively; for the computetime breakdown, we additionally sweep the micro-batch size over 1 and 2. These workloads exercise different compute– communication balances: the smaller model has less compute per communicated byte, while the larger model provides more opportunity to hide host-side work behind GPU execution. In these end-to-end experiments, HOCCL accelerates PP send/recv and DP AllGather. Communication primitives that require reduction semantics, such as ReduceScatter, fall back to NCCL. 10
NCCL-1sm
HOCCL
NCCL-CE
HOCCL
NCCL-1sm
NCCL-4sm
NCCL-1sm
NCCL-2sm
NCCL-2sm
NCCL-4sm
NCCL-4sm
NCCL-8sm
NCCL-8sm
NCCL-4sm
NCCL-8sm
NCCL-8sm
200 150 100 50 0 2 4 8 16
64 256 1024
50 40 30 20 10
2 4 8 16
64 256 1024
200 150 100 50 0
Communication Size (MB)
Communication Size (MB)
64 256 1024
Communication Size (MB)
(b) Config 2: AllGather
(a) Config 1: AllGather
2 4 8 16
Algorithm Bandwidth (GB/s)
HOCCL NCCL-2sm
Algorithm Bandwidth (GB/s)
NCCL-CE NCCL-2sm
Algorithm Bandwidth (GB/s)
Algorithm Bandwidth (GB/s)
HOCCL NCCL-1sm
(c) Config 1: All-to-All
60 40 20 0 2 4 8 16
64 256 1024
Communication Size (MB)
(d) Config 2: All-to-All
Figure 8: Collective communication bandwidth.
HOCCL
HOCCL
10 8 6 2 4 8 16 64 256 1024
transfers rather than providing its own SM-free intra-node path. For inter-node P2P, NCCL’s launch cost remains nearly flat at 8–20 µs, while HOCCL grows from 11.8 µs at 2 MB to 624.6 µs at 1 GB. This growth is expected: HOCCL is hostdriven and submits more DMA/RDMA work as a transfer is divided into more chunks. VCCL has much higher overhead, from 563 µs to 5.4 ms, because it performs blocking RDMA buffer registration whose cost grows with buffer size. This overhead result also identifies a practical limitation of user-buffer registration. Under NIC and GPU driver load, RDMA registration latency is highly variable and can spike to 20–30 ms in our experiments, consistent with prior observations from Meta’s large-scale training deployments [43]. This explains why VCCL can suffer severe jitter despite being host-driven, and motivates HOCCL’s use of reusable staging buffers rather than repeatedly registering application tensors.
NCCL
Inter-node VCCL CPU Launch Time (ms)
CPU Launch Time (µs)
Intra-node NCCL&VCCL
4 2 0 2 4 8 16
64 256 1024
Communication Size (MB)
Communication Size (MB)
(a) Intra-node P2P
(b) Inter-node P2P
Figure 9: CPU-side launch overhead. to-All exchanges. On 8 GPUs (Figure 8c), HOCCL reaches 191.4 GB/s at 1 GB, comparable to NCCL-CE at 187.6 GB/s, and exceeds NCCL-8SM by 24% (154.0 GB/s), while also outperforming all SM-constrained NCCL configurations. At 16 GPUs, bandwidth decreases for both systems because the communication degree and induced intra-node traffic increase. HOCCL reaches 57.7 GB/s at 1 GB, matching NCCL-4SM (57.5 GB/s) and substantially outperforming NCCL-1SM and NCCL-2SM. NCCL-8SM remains faster at 67.7 GB/s, a 17% advantage, because it can spend more SM resources to drive more aggressive communication progress. The scalability trend is therefore consistent across collectives: HOCCL is strongest when SMs are scarce or valuable, while NCCL’s peak bandwidth can still improve by dedicating more SMs to communication.
6.3
End-to-End Training
We next evaluate whether eliminating communication-related SM usage improves full training runs. We integrate HOCCL into Megatron-LM and compare it with NCCL and VCCL on two workloads with different compute–communication balance. We report throughput, per-step compute time, and CPU slack to determine whether HOCCL’s host-driven control path affects the critical path of training. Training throughput. HOCCL improves end-to-end throughput because the SMs saved from communication become available to training kernels. Figures 11a and 11b report per-step throughput for the 6.7B and 32B models. On the 6.7B workload, NCCL achieves 209.8 TFLOP/s on average, VCCL reaches 190.9 TFLOP/s, and HOCCL reaches 214.8 TFLOP/s. HOCCL improves over NCCL by 2.4%, whereas VCCL is 9.0% slower than NCCL. VCCL also exhibits large jitter, ranging from 171.9 to 209.9 TFLOP/s, while HOCCL remains between 212.0 and 218.1 TFLOP/s. The larger 32B workload amplifies HOCCL’s ben-
CPU launch overhead. HOCCL increases CPU launch cost, but the magnitude depends on whether the path is intranode or inter-node. Figures 9a and 9b report CPU-side launch overhead for P2P operations. For intra-node P2P, HOCCL costs 9–12 µs, compared with 5–8 µs for NCCL. The difference comes from HOCCL’s CPU-side sender/receiver coordination for direct DMA writes. VCCL is similar to NCCL in this setting because it falls back to NCCL for intra-node 11
NCCL
time = B0
computeA launch
CPU thread
communicationB launch
GPU stream previous work
CPU slack (ms)
time = A0
CPU slack A = A1-A0>0 CPU slack B = B1-B0≈0
other work
excute communicationB
excute computeA time = A1
CPU VCCL slack
HOCCL
40 20 0 1 2 3 4 5 6 7 8 9 1011121314151617181920 Timestamp
time = B1
(a) The definition of CPU slack
(b) CPU slack over time
Figure 10: CPU slack in asynchronous GPU execution.
Throughput (TFLOP/s)
NCCL
VCCL
HOCCL
NCCL
VCCL
HOCCL
500
220 200
450 180 160
400 200 202 204 206 208 210
200 202 204 206 208 210
Training Step
Training Step
(a) 6.7B model
Extra CPU cost analysis. We next analyze whether the additional CPU time overhead introduced by HOCCL can make the CPU side a training bottleneck, and explain why VCCL exhibits performance fluctuations. Modern training stacks enqueue work asynchronously: while the GPU executes previously issued kernels, the CPU prepares future work. As shown in Figure 10a, we define CPU slack as the difference between an operator’s GPU start time and the end of its CPU launch (i.e., CPU slack = GPU start time - CPU launch end time). CPU slack A is defined as the gap between the CPU submission time of compute kernel A and its actual execution time on the GPU; a positive value indicates that the CPU is not the bottleneck. In contrast, CPU slack B equals zero, meaning that CPU submission is no longer ahead of the execution of prior GPU work. In this case, excessive CPU overhead becomes a system bottleneck. Using torch.profiler, we sample operators and record their CPU enqueue timestamps and GPU start times.
(b) 32B model
Figure 11: End-to-end training throughput Forward pass mbs=1 mbs=2
Backward pass mbs=1 mbs=2
14 12
20
10 15
8 6
10 NCCL
VCCL
HOCCL
(a) Forward compute time
NCCL
VCCL
micro-batch sizes. For the forward pass, HOCCL reduces average compute time from 7.66 ms to 7.55 ms at mbs=1 (1.5%) and from 13.12 ms to 12.41 ms at mbs=2 (5.4%). For the backward pass, HOCCL reduces time from 10.89 ms to 10.19 ms at mbs=1 (6.4%) and from 22.94 ms to 21.08 ms at mbs=2 (8.1%). These reductions directly connect the endto-end throughput gains to the design goal of freeing SM resources for training kernels.
HOCCL
(b) Backward compute time
Figure 12: Compute time per micro-batch efit. NCCL averages 467.8 TFLOP/s, VCCL averages 478.7 TFLOP/s, and HOCCL averages 491.7 TFLOP/s. HOCCL therefore improves over NCCL by 5.1% and over VCCL by 2.7%. VCCL improves over NCCL on average in this workload, but its throughput still varies from 461.6 to 491.4 TFLOP/s because registration and host-side latency occasionally enter the critical path. HOCCL remains both faster and more stable, ranging only from 486.8 to 495.8 TFLOP/s. These results support the central claim that HOCCL converts saved communication SMs into higher training throughput while avoiding VCCL’s registration-induced instability. Impact on compute time. Removing communication kernels from SMs measurably shortens overlapped compute. Figure 12 reports forward and backward compute time for two
Figure 10b shows that NCCL and HOCCL maintain positive slack throughout the measured window, typically 20– 50 ms. Thus, although HOCCL launches more host-side DMA/RDMA work than NCCL, the extra work is absorbed by existing CPU slack and does not shift the bottleneck from GPU execution to CPU scheduling. VCCL has much higher variance and reaches zero slack at timestamp 11, indicating that blocking registration can stall GPU execution. This explains the throughput jitter observed in Figure 11a and Figure 11b, and highlights a trade-off: HOCCL pays predictable host launch overhead to avoid VCCL’s unpredictable registration stalls. 12
7
Conclusion
[8] Ben Cottier, Robi Rahman, Loredana Fattorini, Nestor Maslej, Tamay Besiroglu, and David Owen. The rising costs of training frontier ai models. arXiv preprint arXiv:2405.21015, 2024.
This paper revisits a basic assumption in distributed training systems: collective communication is often optimized for bandwidth, while its consumption of GPU compute resources is overlooked. We present HOCCL to show that communication progress can be moved off GPU SMs and onto dedicated hardware, reducing interference between communication and computation. More broadly, HOCCL suggests that communication design should aim not only for high standalone speed, but also for preserving GPU resources for training. This perspective becomes increasingly important as modern training workloads grow larger and place ever greater pressure on limited GPU compute capacity. It also indicates that future communication libraries should be evaluated not only by their raw communication throughput, but also by how effectively they coexist with end-to-end training execution. Our results show that this perspective is both practical and beneficial, highlighting a promising direction for future large-scale training systems.
[9] DeepSeek-AI. Deepseek-v2: A strong, economical, and efficient mixture-of-experts language model, 2024. [10] Paul Elvinger, Foteini Strati, Natalie Enright Jerger, and Ana Klimovic. Understanding gpu resource interference one level deeper. In ACM Symposium on Cloud Computing (SoCC), 2025. [11] William Fedus, Barret Zoph, and Noam Shazeer. Switch transformers: Scaling to trillion parameter models with simple and efficient sparsity. Journal of Machine Learning Research, 23(120):1–39, 2022. [12] Jordan Hoffmann, Sebastian Borgeaud, Arthur Mensch, Elena Buchatskaya, Trevor Cai, Eliza Rutherford, Diego de Las Casas, Lisa Anne Hendricks, Johannes Welbl, Aidan Clark, Tom Hennigan, Eric Noland, Katie Millican, George van den Driessche, Bogdan Damoc, Aurelia Guy, Simon Osindero, Karen Simonyan, Erich Elsen, Jack W. Rae, Oriol Vinyals, and Laurent Sifre. Training compute-optimal large language models. In Advances in Neural Information Processing Systems 35 (NeurIPS 2022), 2022.
References [1] Josh Achiam, Steven Adler, Sandhini Agarwal, Lama Ahmad, Ilge Akkaya, Florencia Leoni Aleman, Diogo Almeida, Janko Altenschmidt, Sam Altman, Shyamal Anadkat, et al. Gpt-4 technical report. arXiv preprint arXiv:2303.08774, 2023.
[13] Yanping Huang, Youlong Cheng, Ankur Bapna, Orhan Firat, Dehao Chen, Mia Chen, HyoukJoong Lee, Jiquan Ngiam, Quoc V Le, Yonghui Wu, et al. Gpipe: Efficient training of giant neural networks using pipeline parallelism. Advances in neural information processing systems, 32, 2019.
[2] Osayamen Jonathan Aimuyo, Byungsoo Oh, and Rachee Singh. Flashmoe: Fast distributed moe in a single kernel. arXiv preprint arXiv:2506.04667, 2025. [3] Rohan Anil et al. Palm 2 technical report. arXiv preprint arXiv:2305.10403, 2023. [4] Joshua Bakita and James H. Anderson. Hardware Compute Partitioning on NVIDIA GPUs for Composable Systems. In Renato Mancuso, editor, 37th Euromicro Conference on Real-Time Systems (ECRTS 2025), volume 335 of Leibniz International Proceedings in Informatics (LIPIcs), pages 21:1–21:25, Dagstuhl, Germany, 2025. Schloss Dagstuhl – Leibniz-Zentrum für Informatik.
[14] Zixiao Huang, Junhao Hu, Hao Lin, Chunyang Zhu, Yueran Tang, Quanlu Zhang, Zhen Guo, Zhenhua Li, Shengen Yan, Zhenhua Zhu, et al. Stalloc: Enhancing memory efficiency in large-scale model training with spatio-temporal planning. arXiv preprint arXiv:2507.16274, 2025. [15] Changho Hwang, Wei Cui, Yifan Xiong, Ziyue Yang, Ze Liu, Han Hu, Zilong Wang, Rafael Salas, Jithin Jose, Prabhat Ram, Joe Chau, Peng Cheng, Fan Yang, Mao Yang, and Yongqiang Xiong. Tutel: Adaptive mixtureof-experts at scale. Proceedings of Machine Learning and Systems, 5:269–287, 2023.
[5] Garrett Birkhoff. Tres observaciones sobre el algebra lineal. Universidad Nacional de Tucumán. Revista A, 5:147–151, 1946. [6] Mark Chen et al. Evaluating large language models trained on code. arXiv preprint arXiv:2107.03374, 2021.
[16] Ziheng Jiang, Haibin Lin, Yinmin Zhong, Qi Huang, Yangrui Chen, Zhi Zhang, Yanghua Peng, Xiang Li, Cong Xie, Shibiao Nong, Yulu Jia, Sun He, Hongmin Chen, Zhihao Bai, Qi Hou, Shipeng Yan, Ding Zhou, Yiyao Sheng, Zhuo Jiang, Haohan Xu, Haoran Wei, Zhang Zhang, Pengfei Nie, Leqi Zou, Sida Zhao, Liang Xiang, Zherui Liu, Zhe Li, Xiaoying Jia, Jianxi Ye, Xin
[7] Ziteng Chen, Xiaohe Hu, Menghao Zhang, Yanmin Jia, Yan Zhang, Mingjun Zhang, Da Liu, Fangzheng Jiao, Jun Chen, He Liu, et al. An efficient, reliable and observable collective communication library in large-scale gpu training clusters. arXiv preprint arXiv:2510.00991, 2025. 13
Jin, and Xin Liu. Megascale: Scaling large language model training to more than 10,000 gpus. In 21st USENIX Symposium on Networked Systems Design and Implementation (NSDI 24), pages 745–760, 2024.
[27] NVIDIA. Fusing Communication and Compute with New Device API and Copy Engine Collectives in NVIDIA NCCL 2.28. NVIDIA Technical Blog, 2025. Accessed: 2026-04-19.
[17] Dmitry Lepikhin, HyoukJoong Lee, Yuanzhong Xu, Dehao Chen, Orhan Firat, Yanping Huang, Maxim Krikun, Noam Shazeer, and Zhifeng Chen. Gshard: Scaling giant models with conditional computation and automatic sharding. In International Conference on Learning Representations, 2021.
[28] NVIDIA. Collective Operations – NCCL 2.29.7 Documentation. NVIDIA, 2026. Accessed 2026-04-20. [29] NVIDIA. CUDA Programming Guide. NVIDIA, March 2026. Release 13.2, accessed 2026-04-20. [30] NVIDIA. CUDA Stream Semantics. NVIDIA, 2026. NCCL 2.29.7 documentation, accessed 2026-04-20.
[18] Shen Li, Yanli Zhao, Rohan Varma, Omkar Salpekar, Pieter Noordhuis, Teng Li, Adam Paszke, Jeff Smith, Brian Vaughan, Pritam Damania, and Soumith Chintala. Pytorch distributed: Experiences on accelerating data parallel training. Proceedings of the VLDB Endowment, 13(12):3005–3018, 2020.
[31] NVIDIA. NCCL Tests. https://github.com/ NVIDIA/nccl-tests, 2026. GitHub repository, accessed 2026-04-21. [32] NVIDIA. NVIDIA Collective Communications Library (NCCL). NVIDIA, 2026. Version 2.29.7, NVIDIA documentation, accessed March 26, 2026.
[19] Shigang Li and Torsten Hoefler. Chimera: Efficiently training large-scale neural networks with bidirectional pipelines. In Proceedings of the International Conference for High Performance Computing, Networking, Storage and Analysis, SC ’21, 2021.
[33] NVIDIA. NVIDIA GB200 Developer Kit Specifications. NVIDIA, 2026. Official product specifications, accessed: 2026-03-30. [34] NVIDIA Developer Blog. CUTLASS: Fast Linear Algebra in CUDA C++. NVIDIA Developer Blog, December 2017.
[20] Yujia Li, David Choi, Junyoung Chung, Nate Kushman, Julian Schrittwieser, Rémi Leblond, Tom Eccles, James Keeling, Felix Gimeno, Agustin Dal Lago, et al. Competition-level code generation with alphacode. Science, 378(6624):1092–1097, 2022.
[35] Kishore Punniyamurthy, Khaled Hamidouche, and Bradford M. Beckmann. Optimizing distributed ml communication with fused computation-collective operations. arXiv preprint arXiv:2305.06942, 2023.
[21] Yujia Li et al. Competition-level code generation with alphacode. Science, 378(6624):1092–1097, 2022.
[36] PyTorch. CUDA semantics. https: //docs.pytorch.org/docs/stable/ notes/cuda.html#memory-management. PyTorch 2.11 documentation, section “Memory management”, accessed: 2026-03-29.
[22] Jiamin Lu, Jingwei Sun, Yunlong Xu, Peng Sun, and Guangzhong Sun. Conco: Optimizing compilation of concurrent tensor programs on shared gpu. In Proceedings of the 2025 International Conference on Supercomputing, ICS ’25, Salt Lake City, UT, USA, 2025. ACM.
[37] PyTorch. PyTorch Profiler. https://docs. pytorch.org/tutorials/recipes/ recipes/profiler_recipe.html. PyTorch Tutorials, v2.11.0+cu130, accessed: 2026-03-29.
[23] Deepak Narayanan, Aaron Harlap, Amar Phanishayee, Vivek Seshadri, Nikhil R Devanur, Gregory R Ganger, Phillip B Gibbons, and Matei Zaharia. Pipedream: Generalized pipeline parallelism for dnn training. In Proceedings of the 27th ACM symposium on operating systems principles, pages 1–15, 2019.
[38] PyTorch Contributors. Distributed communication package - torch.distributed. PyTorch, 2026. Accessed 202604-20.
[24] NVIDIA. Megatron-lm. https://github.com/ NVIDIA/Megatron-LM. GitHub repository, accessed: 2026-03-30.
[39] Penghui Qi, Xinyi Wan, Guangxing Huang, and Min Lin. Zero bubble (almost) pipeline parallelism. In International Conference on Learning Representations, 2024.
[25] NVIDIA. NVIDIA A100 Tensor Core GPU Architecture. NVIDIA, 2020. White paper.
[40] Ruoyu Qin, Zheming Li, Weiran He, Jialei Cui, Feng Ren, Mingxing Zhang, Yongwei Wu, Weimin Zheng, and Xinran Xu. Mooncake: Trading more storage for less computation—a kvcache-centric architecture for
[26] NVIDIA. RDMA Aware Networks Programming User Manual: Key Concepts. NVIDIA, 2024. Accessed: 2026-04-03. 14
serving llm chatbot. In 23rd USENIX conference on file and storage technologies (FAST 25), pages 155–170, 2025. [41] Samyam Rajbhandari, Jeff Rasley, Olatunji Ruwase, and Yuxiong He. Zero: Memory optimizations toward training trillion parameter models. In Proceedings of the International Conference for High Performance Computing, Networking, Storage and Analysis, SC ’20, pages 20:1–20:16, 2020. [42] Mohammad Shoeybi, Mostofa Patwary, Raul Puri, Patrick LeGresley, Jared Casper, and Bryan Catanzaro. Megatron-lm: Training multi-billion parameter language models using model parallelism. arXiv preprint arXiv:1909.08053, 2019. [43] Min Si, Pavan Balaji, Yongzhou Chen, Ching-Hsiang Chu, Adi Gangidi, Saif Hasan, Subodh Iyengar, Dan Johnson, Bingzhe Liu, Regina Ren, et al. Collective communication for 100k+ gpus. arXiv preprint arXiv:2510.20171, 2025. [44] Jaeyong Song, Jinkyu Yim, Jaewon Jung, Hongsun Jang, Hyung-Jin Kim, Youngsok Kim, and Jinho Lee. Optimus-cc: Efficient large nlp model training with 3d parallelism aware communication compression. In Proceedings of the 28th ACM International Conference on Architectural Support for Programming Languages and Operating Systems, Volume 2, 2023. [45] The Open Group. <semaphore.h> — Semaphores. The Open Group, 2018. POSIX.1-2017, IEEE Std 1003.1-2017, accessed 2026-04-20. [46] Chenggang Zhao, Shangyan Zhou, Liyue Zhang, Chengqi Deng, Zhean Xu, Yuxuan Liu, Kuai Yu, Jiashi Li, and Liang Zhao. Deepep: an efficient expertparallel communication library. https://github. com/deepseek-ai/DeepEP, 2025. [47] Yanli Zhao, Andrew Gu, Rohan Varma, Liang Luo, Chien-Chin Huang, Min Xu, Less Wright, Hamid Shojanazeri, Myle Ott, Sam Shleifer, Alban Desmaison, Can Balioglu, Pritam Damania, Bernard Nguyen, Geeta Chauhan, Yuchen Hao, Ajit Mathews, and Shen Li. Pytorch fsdp: Experiences on scaling fully sharded data parallel. Proceedings of the VLDB Endowment, 16(12):3848–3860, 2023. [48] Size Zheng, Jin Fang, Xuegui Zheng, Qi Hou, Wenlei Bao, Ningxin Zheng, Ziheng Jiang, Dongyang Wang, Jianxi Ye, Haibin Lin, Li-Wen Chang, and Xin Liu. Tilelink: Generating efficient compute-communication overlapping kernels using tile-centric primitives. arXiv preprint arXiv:2503.20313, 2025.
15