Full text not available for this paper

Summary (Overview)

  • Cobalt is a novel distributed training framework for Mixture-of-Experts (MoE) models that leverages the observation that many expert pairs are frequently co-activated by the same tokens.
  • It introduces a two-stage expert layout planner that (1) periodically performs a global re-layout of experts based on co-activation and workload statistics, and (2) performs per-step intra-node replica rebalancing to ensure workload balance.
  • A communication-aware task assignment algorithm routes tokens to the fewest remote nodes possible, significantly reducing cross-node traffic.
  • Experiments on 32 NVIDIA B200 GPUs show Cobalt achieves up to 1.53–2.41× speedup (1.28–1.89× on average) over existing frameworks, while reducing cross-node token transfer volume by 75.74%–99.26%.
  • This is the first work to exploit expert co-activation for communication reduction and workload balance in distributed MoE training.

Introduction and Theoretical Foundation

Background & Motivation: MoE architectures scale LLMs by routing each token to a small subset of experts, keeping per-token computation nearly constant. Training large MoE models relies on Expert Parallelism (EP), which distributes expert replicas across GPUs and uses all-to-all communication. EP is bottlenecked by two factors:

  1. Cross-node token transfers: Inter-node bandwidth (e.g., 400 GB/s InfiniBand) is significantly lower than intra-node bandwidth (e.g., 1.8 TB/s NVLink). Communication can account for 58%–83% of training time.
  2. Skewed expert workloads: Uneven token-to-expert assignments lead to imbalanced computation across GPUs.

Key Observation (Expert Co-activation): The authors empirically track the set of active tokens for each expert and compute the Ochiai similarity between expert pairs:

Ochiai(A,B)=∣A∩B∣∣A∣⋅∣B∣\text{Ochiai}(A, B) = \frac{|A \cap B|}{\sqrt{|A| \cdot |B|}}

where AA and BB are the sets of tokens activating experts AA and BB. During training of GLM-4.5-Air, the average Ochiai similarity rises from 0.41 to 0.89 (Layer 35) and from 0.28 to 0.86 (Layer 41). This indicates that many expert pairs are co-activated by the same tokens, and this co-activation strengthens over training. Critically, when two selected experts reside on the same remote node, node-level token deduplication means the token is only sent cross-node once, so co-locating co-activated experts can dramatically reduce cross-node traffic.

Problem Formulation: The paper formalizes the joint optimization of expert layout and task assignment. The communication cost is modeled as the total number of cross-node token transfers:

V=∑t=1T∣⋃e∈Dt{N(rte)}∖{N(st)}∣V = \sum_{t=1}^{T} \left| \bigcup_{e \in D_t} \{N(r_t^e)\} \setminus \{N(s_t)\} \right|

where TT is the number of tokens, N(r)N(r) is the node containing rank rr, DtD_t is the set of selected experts for token tt, sts_t is the source rank, and rter_t^e is the rank assigned to token-expert task (t,e)(t,e).

The full optimization problem is:

arg⁡min⁡{Pr},{rte}V\arg\min_{\{P_r\}, \{r_t^e\}} V s.t.∑rI[e∈Pr]≥1,∀e,∣Pr∣≤(1+ρ)E/G,∀r\text{s.t.} \sum_r \mathbb{I}[e \in P_r] \geq 1, \forall e, \quad |P_r| \leq (1+\rho)E/G, \forall r e∈Prte,∀t,e∈Dt,∑e∈Prwre≤(1+ϵ)KT/G,∀re \in P_{r_t^e}, \forall t, e \in D_t, \quad \sum_{e \in P_r} w_r^e \leq (1+\epsilon)KT/G, \forall r

Here, PrP_r is the set of experts with a replica on rank rr, ρ\rho is the replica budget, ϵ\epsilon is the workload imbalance tolerance, EE is the number of experts, GG is the number of ranks, and KK is the number of selected experts per token.

Methodology

System Architecture: Cobalt separates expert layout planning from task assignment.

1. Two-Stage Expert Layout Planner

Stage 1: Periodic Global Planning. Every 10 training steps, a global re-layout is computed using historical statistics:

  • w^e\hat{w}_e: EMA of expert ee's assigned computation task count
  • c^ij\hat{c}_{ij}: EMA of the number of tokens selecting both experts ii and jj

The planning problem is reformulated to maximize total co-activation under workload constraints:

arg⁡max⁡{Pr},{wer′}∑n∑i<jc^ijI[i∈P(n)∧j∈P(n)]\arg\max_{\{P_r\}, \{w'_{er}\}} \sum_n \sum_{i<j} \hat{c}_{ij} \mathbb{I}[i \in P(n) \wedge j \in P(n)]

This is re-expressed as a Mixed-Integer Linear Programming (MILP) problem (Eq. 4 in the paper) and solved using libraries like SCIP. The solution yields the expert layout {Pr}\{P_r\} and expected workload division {wer′}\{w'_{er}\}.

Stage 2: Per-Step Intra-Node Rebalancing. Between global updates, Cobalt adjusts replica placement within each node at every training step to rebalance workloads. This is formulated as a variant of the min-sum set partitioning problem, solved via a greedy bin-packing approach that sorts replicas by EMA workloads and partitions them greedily. This is fast enough to be overlapped with model computation.

2. Communication-Aware Task Assignment

For each token, the algorithm (Algorithm 1) works in two steps:

  • Step 1 (Selecting execution nodes): For experts not on the source node, solve a minimum set-cover problem to find the fewest remote nodes covering all remaining experts. This is done via bitmask enumeration over the N−1N-1 remote nodes.
  • Step 2 (Selecting execution ranks): Uniformly sample an execution rank from the available replicas on the chosen node, ensuring equal expected workload distribution.

The complexity is O(T⋅2N−1)O(T \cdot 2^{N-1}), which is negligible since NN is small (e.g., N=4N=4 for 32 GPUs with 8 GPUs/node).

3. Implementation Details

  • Built on PyTorch with DeepEP as the communication backend for node-level token deduplication
  • MILP and min-sum partitioning solved on CPU; task assignment implemented as a GPU kernel using Triton
  • Global layout updates triggered every 10 steps; intra-node rebalancing performed per-step and overlapped with computation

Empirical Validation / Results

Setup: Experiments on 4 NVIDIA B200 GPU servers (32 GPUs) with 1.8 TB/s NVLink and 400 GB/s InfiniBand. Models evaluated: Hunyuan3 (60B), GLM-4.5-Air (106B), DeepSeek-V3 (140B). Baselines: EP+FSDP, SmartMoE, LAER-MoE (all using DeepEP).

End-to-End Performance (Figure 6):

ModelEP DegreeEP+FSDP (s)SmartMoE (s)LAER-MoE (s)Cobalt (s)Speedup vs EP+FSDP
Hunyuan3EP164.973.92 (↑1.27×)3.87 (↑1.28×)3.62 (↑1.37×)1.37×
Hunyuan3EP3216.8213.52 (↑1.24×)12.26 (↑1.37×)10.97 (↑1.53×)1.53×
GLM-4.5-AirEP1612.308.45 (↑1.46×)7.83 (↑1.57×)7.01 (↑1.76×)1.76×
GLM-4.5-AirEP3230.2525.68 (↑1.18×)19.25 (↑1.57×)12.57 (↑2.41×)2.41×
DeepSeek-V3EP168.587.41 (↑1.16×)5.61 (↑1.53×)4.30 (↑2.00×)2.00×
DeepSeek-V3EP3219.0314.91 (↑1.28×)12.52 (↑1.52×)8.21 (↑2.32×)2.32×

Key Results:

  • Cross-node communication reduction: Cobalt reduces cross-node traffic by 75.74%–99.26% compared to EP+FSDP, and by 30.71–78.29% compared to LAER-MoE.
  • Workload imbalance: Cobalt consistently achieves the lowest workload imbalance ratio (e.g., 1.13–1.24 vs. 3.35–6.66 for EP+FSDP).
  • Ablation Study (GLM-4.5-Air, EP=16): Intra-node rebalancing lowers imbalance ratio from 3.35 → 1.28 (1.22× speedup); periodic global update reduces cross-node communication by 95.3% (1.43× speedup); communication-aware task assignment further reduces communication by 98.2%. Combined: 1.53× speedup.
  • Overhead: Intra-node rebalancing takes negligible time and is overlapped; global planning takes 1.21–2.61 seconds but is overlapped with computation; task assignment overhead is <1 ms.
  • Sensitivity: Performance is robust to varying ρ\rho (0.25–0.75) and ϵ\epsilon (0.10–0.30), with <5% variation in training time.
  • Scalability: Increasing GPUs from 8 to 32 improves throughput by 3.20–3.27×.

Theoretical and Practical Implications

Theoretical Contributions:

  • Introduces a new perspective on MoE training optimization: leveraging pairwise expert co-activation rather than per-expert workload statistics alone
  • Provides a formal optimization framework (Eq. 2) that jointly considers expert layout, task assignment, communication cost, and workload balance
  • Demonstrates that co-activation patterns strengthen over training (Ochiai similarity increases from ~0.3 to ~0.9), suggesting structured specialization in MoE models

Practical Implications:

  • For MoE training systems: Cobalt shows that communication-aware expert placement can dramatically reduce cross-node traffic (up to 99%), which is critical given the bandwidth hierarchy in modern GPU clusters
  • For hardware design: The results highlight the importance of high intra-node bandwidth and node-level communication deduplication
  • For model design: The observed co-activation phenomenon may inform future MoE architecture design, potentially encouraging gating mechanisms that explicitly promote expert co-activation

Conclusion

Main Takeaways:

  • Cobalt is the first framework to exploit expert co-activation for both communication reduction and workload balance in distributed MoE training
  • The two-stage planner (periodic global + per-step intra-node) effectively adapts to evolving co-activation patterns while keeping overhead low
  • Communication-aware task assignment with minimum set-cover optimization significantly reduces cross-node traffic
  • Achieves state-of-the-art performance: up to 2.41× speedup and 99.26% reduction in cross-node communication

Future Directions:

  • Refining the optimization problem solving (current approach uses EMA statistics and heuristics, which may be sub-optimal)
  • Evaluating Cobalt on larger clusters with other parallelism strategies (e.g., combining with data/sequence parallelism)
  • Exploring how the co-activation phenomenon might inform future MoE architecture and gating design

Related papers