Conceptio › Archive › arXiv CS
arXiv CSopen access

TopoEP: Topology-Aware Load Balancing for Expert-Parallel MoE Training

Jiacheng Zhu et al. · arxiv_cs
arXiv CS · Papers · License: Open Access
Open Source ↗Direct PDF ↓
clouddistributed-computingparallel-computing
distributed computing, parallel computing, cloud

TopoEP: Topology-Aware Load Balancing for Expert-Parallel MoE Training Jiacheng Zhu, Xie Zhao, Gongming Zhao, Hongli Xu, Yao Fei, and Jin Fang University of Science and Technology of China {zhu_jc,zhaoxie,yao_fei,fangjin98}@mail.ustc.edu.cn

arXiv:2609.35481v1 [cs.DC] 28 Sep 2026

{gmzhao,xuhongli}@ustc.edu.cn

Abstract

As EP groups grow, rank-level load imbalance can become a major system bottleneck [12]. Top-k gating may concentrate token–expert assignments on a small set of hot experts whose identities and loads vary across layers and microbatches [20, 28, 34]. Under static expert placement, this expert-level skew can translate into uneven aggregate loads across ranks, making ranks that host hot experts stragglers. Because an MoE layer completes only after every participating rank finishes, even one overloaded rank can prolong the entire EP stage. The same skew can also produce imbalanced dispatch / combine traffic, creating communication hotspots and further reducing end-to-end training efficiency [21, 33, 38]. Existing approaches address this imbalance at either the model or system level. Model-level approaches modify routing through auxiliary losses, routing biases or capacityinduced token dropping, potentially affecting model quality [6, 13, 33]. System-level methods preserve each token’s logical expert assignments while replicating hot experts and redistributing their assigned tokens across expert instances [5, 10, 19]. Such methods need to generate each plan rapidly because planning lies on the critical path before token dispatch. Existing online approaches commonly rely on hostside optimization solvers, whose GPU–CPU synchronization, CPU solving and plan distribution add directly to this criticalpath latency. More recent systems such as UltraEP [34] and MoonEP [1] adopt GPU-native load balancing but target a single scale-up domain, overlooking the heterogeneous costs of intra-node NVLink and inter-node RDMA. Consequently, when extended to scale-out, topology-unaware planning may reduce rank-load imbalance while incurring costly crossdomain token, expert-parameter, and replica-gradient transfers, offsetting the performance gains from improved load balance. These limitations motivate an online system-level load balancer that improves training efficiency without altering Top-k expert selection or compromising model quality. Realizing such a system presents three coupled challenges: planning quality, solver latency and execution overhead. First, the planner needs to derive expert-replication and tokenrerouting decisions from the current Top-k gating result under

Dynamic routing creates severe load imbalance in large-scale expert-parallel Mixture-of-Experts (MoE) training, turning GPUs that host hot experts into stragglers. As each MoE layer waits for its slowest rank, these stragglers prolong the expertparallel stage and reduce overall training efficiency. Existing expert-parallelism load-balancing (EPLB) systems commonly compute load-balancing plans on the CPU, incurring device– host data transfers and cross-rank synchronization that make scheduling at every layer and microbatch expensive. Their planning formulations also overlook the hierarchical communication costs of modern scale-up and scale-out GPU clusters. We present TopoEP, a GPU-native, topology-aware loadbalancing system for large-scale MoE training. At each MoE layer and training microbatch, TopoEP converts the current routing result into hot-expert replication and tokenrerouting decisions and executes the resulting plan without data-dependent host synchronization, reducing critical-path overhead. To generate these decisions, TopoEP uses a deterministic GPU solver that performs inter-node placement followed by intra-node refinement, allowing all ranks to independently produce bitwise-identical plans. On a 32-GPU NVIDIA H800 cluster, integrating TopoEP with MegatronLM improves end-to-end training throughput by 6.2%–11.4% across three representative MoE models.

1

Introduction

Mixture-of-Experts (MoE) is a widely adopted architecture for scaling large language models (LLMs). By activating only a small subset of experts for each token, MoE increases model capacity without proportionally increasing per-token computation [2, 4, 6, 13, 26]. Expert parallelism (EP) enables largescale MoE training by partitioning experts across GPUs. The dispatch phase sends token activations to the ranks hosting the router-selected experts, and the combine phase returns their outputs to the source ranks [7, 11, 13, 25]. EP has therefore become a standard strategy for distributing expert parameters and computation across ranks in large-scale MoE systems. 1

a limited replica budget. A high-quality plan reduces the maximum expert-computation time across ranks while accounting for communication costs. Second, plan generation lies on the critical path between Top-k gating and token dispatch, leaving little opportunity to hide its latency. At large EP scales, rapidly solving this joint replication-and-routing problem across many experts and ranks is challenging because any solver delay directly postpones token dispatch and can offset the gains from improved load balance. Finally, applying each plan introduces expert-parameter transfers, replica-gradient aggregation and additional memory usage. These overheads need to remain low enough for improved load balance to translate into end-to-end performance gains. We therefore present TopoEP, a GPU-native load-balancing system that transforms the routing outcome of each MoE layer and microbatch into an expert-replication and tokenrerouting plan. For high-quality planning, TopoEP jointly determines replica placement and token allocation among expert instances, subject to a per-rank replica limit. Inter-node placement reduces RDMA traffic, while intra-node refinement mitigates residual rank-load imbalance. For low-latency solving, TopoEP exploits block-level concurrency, shared-memory caching and warp-level reductions to execute this deterministic procedure entirely on the GPUs. Given the same global routing matrix and configuration, every rank independently produces a bitwise-identical plan without data-dependent host coordination or an additional plan broadcast. For lowoverhead execution, TopoEP uses reusable replica buffers and device-initiated communication through TMA [22] and NCCL GIN [9]. It orchestrates the compute and communication streams through a two-chunk pipeline to reduce exposed communication overhead. We integrate TopoEP into the Megatron-LM training framework and evaluate it on a 32-GPU NVIDIA H800 cluster. Compared with the Megatron-LM baseline, TopoEP achieves lower rank-load imbalance, reduces cross-domain token–expert assignments by 96.2%–98.8%, and improves end-to-end training throughput by 6.2%–11.4% across three representative MoE models, without observable degradation in training-loss convergence. In summary, this paper makes the following contributions:

reducing the critical-path overhead of parameter and gradient transfers as well as token communication. • End-to-end system and evaluation. We integrate TopoEP into Megatron-LM and evaluate solver scalability, load balance, communication locality, end-toend throughput, and training stability on a cluster of 32 NVIDIA H800 GPUs.

2

Background and Motivation

2.1

Expert-Parallel MoE Execution

A dense Transformer layer applies the same feed-forward network (FFN) to every token [32]. In contrast, an MoE layer contains multiple independently parameterized FFNs, called experts. For each token, a trainable router selects the Topk experts based on its hidden representation. By activating only a small subset of experts, MoE increases model capacity without a proportional increase in computation [6, 13, 26]. As the number and size of experts increase, a single GPU can no longer store and execute every expert in an MoE layer. Large-scale MoE models therefore use expert parallelism (EP) [13], which partitions experts across multiple GPUs. A token activation originating on one GPU may consequently be routed to an expert hosted on another, distributing both expert parameters and computation across devices [13, 25]. Expert 0 Rank 0

Router

Σ

Y₀

Expert 1

Dispatch

Combine Expert 2

Rank 1

Σ

Router

Y₁

Expert 3 tokens from Rank 0

tokens from Rank 1

Figure 1: Execution flow of an expert-parallel MoE layer. As shown in Figure 1, each rank first applies the router to its local tokens. An all-to-all dispatch sends the resulting token activations to the ranks hosting the selected experts, where the experts evaluate their assigned tokens [11]. Since expert computation is dominated by matrix multiplications, its cost generally grows with the number of assigned tokens [8, 37]. A reverse all-to-all then returns the expert outputs to their source ranks, where they are combined using the corresponding routing weights. The realized expert load becomes available only after routing. An online planner that uses this load must therefore make the resulting plan available to all participating ranks before dispatch begins. Any exposed planning latency directly delays

• Topology-aware deterministic GPU planning. We formulate joint replica-placement and token-routing decisions across scale-up and scale-out fabrics and develop a two-stage GPU algorithm that trades cross-domain token traffic against replica-transfer cost before refining rank loads within each domain. Its deterministic parallel execution lets every rank independently generate an identical plan. • GPU-resident plan execution. TopoEP applies each plan using reusable replica buffers, device-initiated TMA / GIN communication, and a two-chunk pipeline, 2

dispatch and expert computation, placing the planner on the critical path. This timing constraint is central to Section 2.4.

2.2

ties change across adjacent microbatches: averaged across layers and microbatch pairs, 15.9% and 20.0% of the eight highest-load experts change on DAPO-Math and StarCoderData, respectively. Together, these results show that expert hotspots are common and vary across both layers and microbatches. Consequently, a fixed expert placement can become stale as routing loads change [10, 11]. Under a fixed expert-to-rank mapping, expert-level skew can translate into rank-level imbalance. Ranks hosting hot experts process more tokens and become stragglers, while other ranks finish earlier and wait at synchronization points [10,18]. The same skew also creates uneven communication [14, 16]: hot ranks receive more token activations during dispatch and return more outputs during combine. Consequently, they can become stragglers in both computation and communication.

Dynamic Expert Load Imbalance

1

7

7

13

13

19 25 31 37

0.04

19 25

0.02

31 37

43

43 0

16 32 48 64 80 96 112 Logical Expert ID (a) DAPO-Math DAPO-Math

10

5

0

10

20 30 MoE layer

40

(c) Layer-wise expert load imbalance

50

Topology-Dependent Cost

Equal maximum rank loads do not imply equal communication costs. Scale-out EP groups span NVLink domains interconnected by RDMA, so replica placement determines which token, parameter, and replica-gradient transfers cross domains.

0.00 0

16 32 48 64 80 96 112 Logical Expert ID (b) StarCoderData

N tokens for E

Node 0

StarCoderData Hot-expert turnover (%)

Expert-load imbalance (max/mean)

2.3 Assignment share

1

MoE Layer

MoE Layer

Routing decisions depend on each token’s hidden representation, so expert loads can be skewed and vary with the input [17, 39]. When a batch contains many tokens associated with a small set of learned patterns, the corresponding experts may be selected disproportionately often. These hot experts receive substantially more assignments than the mean, while other experts receive few or none. We profile all 48 MoE layers of Qwen3-30B-A3B [31] on two datasets: DAPOMath [35] for mathematical reasoning and StarCoderData [15] for code. Each layer contains 128 logical experts and routes each token to eight experts.

40%

Rank 0

NVLink/NVSwitch

Rank 1

N tokens 20%

0%

16 24 32 40 Micro-batch index

48

Rank 3

Rank 2

E′ (replica) Load = N

E (main) Load = N

Token traffic: 64 MiB

NIC NIC 8

N tokens for E

Node 1

NVLink/NVSwitch

NIC NIC

(a) Replica on Node 1

(d) Turnover among the eight busiest experts

N tokens for E

Node 0

Node 1

Figure 2: Expert-load skew across the 48 MoE layers of Qwen3-30B-A3B on DAPO-Math and StarCoderData. (a)– (b) Fraction of each layer’s token–expert assignments routed to each expert (horizontal axis: logical expert ID 0–127; vertical axis: MoE layer 1–48). (c) Per-layer maximum-to-mean expert-load ratio. (d) Fraction of the eight highest-load experts that changes between adjacent microbatches.

Rank 1

Rank 0

E′ (replica) Load = N

N tokens NIC NIC

N tokens for E NVLink/NVSwitch

NVLink/NVSwitch

Rank 2

E (main) Load = N

Replica traffic: 27 MiB

Rank 3

NIC NIC

(b) Replica on Node 0

Figure 3: Inter-node communication for two replica placements. Each node contributes N = 4096 tokens to expert E. Both placements use one replica and have the same maximum per-rank load, Lmax = N. (a) A replica on Node 1 incurs 64 MiB of token traffic (solid arrow). (b) A replica on Node 0 incurs 18 MiB of parameter traffic and 9 MiB of replicagradient traffic (dashed arrow).

Figure 2(a)–(b) visualizes the expert-load distribution in every layer. Each row represents one layer, each column one logical expert, and darker cells indicate that the expert receives more tokens in that layer. Within most rows, a small subset of experts receives several times the load expected under uniform routing, while the remaining experts receive substantially less. Figure 2(c) quantifies this skew. The maximum-to-mean ratio averages 7.34 on DAPO-Math and 7.08 on StarCoderData, remains above five in nearly every layer, and peaks at 12.93 and 8.85, respectively. Under synchronized EP execution, such a hotspot can increase its host rank’s workload and delay the entire layer. Figure 2(d) further shows that hotspot identi-

Figure 3 illustrates this tradeoff with two nodes, each forming an NVLink domain. Each node contributes N = 4096 tokens routed to expert E, whose main instance resides on Node 1. Both placements use one replica, E ′ , and assign N tokens to each instance, giving the same maximum per-rank 3

(a) Prior: host-side EPLB planner

load, Lmax = N. The example uses Qwen3-30B-A3B’s 2,048dimensional activations and 9 MiB of parameters per expert, with BF16 activations, parameters, and replica gradients. In Figure 3(a), inter-node token traffic totals 64 MiB over the forward and backward passes. Figure 3(b) keeps token routes local, replacing this traffic with 18 MiB of parameter traffic and 9 MiB of replica-gradient traffic. Two parameter fetches are required because replica buffers are reused across layers. The second placement therefore reduces RDMA traffic by a factor of approximately 2.4 while preserving the maximum per-rank load. This example motivates jointly accounting for rank load, token routing, and replica-transfer cost.

2.4

CPU L5 L2 L3 L1 L4 CPU plan solve (host) …………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………… ……………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………

………………………………………………………………………………………----------------------------------------………………………………………………………… GPU K1 K2 K3 K4 idle bubble K5 (device)………………………………………………………………………………………---------------------------------------- …………………………………………………………

GPU busy

(b) Ours: GPU-native planner

CPU …………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………… L1 L2 L3 L4 L5 (host) …………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………… ~95%

GPU K1 K2 K3 K4 K5 (device)…………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………… ……………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………

GPU busy

time planner on device

time saved planner-input (D2H copy)

launch planner planner-output (H2D copy)

Figure 4: Critical-path execution of host-side and GPU-native EPLB planning. GPU-native planning removes the host round trip between routing and dependent dispatch operations.

Planning on the Critical Path

An online EPLB planner that adapts to the current microbatch depends on the routing decisions produced by Top-k gating. Its replica-placement and token-routing decisions must then be available before the dependent dispatch operations can begin. This dependency places planning on the execution path from gating to dispatch.

GPU Device Routing result

Symmetric buffer Token workspace

World view

Load balancer

Host-side planning. When the planner runs on the CPU, device-resident load statistics must be made available to the host, and the resulting plan must be returned to the GPUs. For a serialized implementation, the latency from device-input readiness to device-plan readiness can be decomposed as Thost-plan = Tsync + TCPU-solve + Treturn ,

~60%

Network Dispatch

(1)

TMA

Intra-node

NCCL GIN

RDMA

Combine

Expert slots

Combine

SMs runtime solver

Token rerouting

NVLink

Dispatch

Expert placement

Expert Transfer

Inter-node

Figure 5: Architecture and execution flow of TopoEP.

where Tsync includes synchronization and the D2H input transfer, TCPU-solve is the CPU solving time, and Treturn is the H2D result transfer. These serialized stages block dependent GPU work until the plan becomes available. As illustrated in Figure 4, host-side planning causes a datadependent GPU–CPU–GPU round trip between gating and dispatch. Although asynchronous kernel submission can overlap host execution with GPU work, dependent dispatch operations must wait for CPU solving and plan transfer to complete. When insufficient independent GPU work is available, this dependency exposes a GPU idle interval. For a fixed workload and GPU configuration, such stalls lengthen the training step and reduce model FLOPs utilization (MFU). These delays can recur at every MoE layer and microbatch, increasing their impact on training efficiency. GPU-native planning removes the host round trip by generating plans on the device for direct consumption by subsequent GPU operations, eliminating this source of exposed control latency.

dients. The net-benefit condition is Tsaved > Tplan + Ttransfer .

(2)

The GPU solver must therefore generate each plan with low latency, while the execution path minimizes or overlaps parameter and gradient transfers. These requirements motivate the GPU-native design of TopoEP, described in Section 3.

3

System Design

Figure 5 summarizes the control and data flows of TopoEP. The globally gathered routing information serves as input to the load balancer, which generates the replica placement and token routing (Section 3.1). The replica placement drives expert-parameter transfers into reusable replica slots, using TMA over intra-node NVLink and NCCL GIN over internode RDMA (Section 3.2). The token-routing decisions drive dispatch and combine through the token workspace, while the two-chunk pipeline overlaps communication with expert computation (Section 3.3). Together, these stages keep both planning and execution on the GPUs.

End-to-end benefit. Dynamic EPLB reduces training time only when the saved expert FFN computation exceeds the additional planning and expert-transfer overhead. Let Tsaved denote the expert-computation time saved relative to static placement, Tplan the GPU solver latency, and Ttransfer the exposed cost of transferring expert parameters and replica gra4

Node 1

Node 0 rank1

rank0 Main-params

E0

E1

E2

E3

E4

E5

S0

S1

S2

E1

E6

S2

TMA

Grad accum.

R0

R1

R2

E6

E7

E8

S0

S1

S2

R0

R1

R2

GIN

TMA

Replica

rank associated with each occupied replica slot and therefore determines both the parameter source and gradient destination. As shown in Figure 6, the CUDA kernels in TopoEP consume this mapping directly: the forward pass transfers expert parameters from main instances to replica slots, and the backward pass returns replica gradients to the gradientaccumulation buffers on the corresponding main ranks. Initiating these transfers directly from the GPU avoids copying the schedule to the CPU and introducing host synchronization between planning and execution. TopoEP selects the device-initiated transfer mechanism according to the network path. Tensor Memory Accelerator (TMA) [22] handles intra-node transfers over NVLink, while NCCL GPU-Initiated Networking (GIN) [9] provides one-sided inter-node RDMA over InfiniBand or RoCE. GIN exposes registered communication buffers through symmetricmemory windows. After the buffers are collectively registered during initialization, a GPU kernel addresses remote memory using a window handle, peer rank, and window-relative offset. Assume that the main instances of the E logical experts are evenly distributed across P EP ranks and that each expert’s parameters occupy |W | bytes. For Nslot replica slots per rank, the three communication buffers have the following per-rank capacities and are reused throughout training:

rank2

GIN

R0

R1

R2

FWD

BWD

Figure 6: Device-initiated replica transfer using intra-node TMA and inter-node NCCL GIN: forward pulls expert parameters into replica slots, and backward returns replica gradients to their main ranks.

3.1

GPU-Native Load Balancer

After Top-k gating, TopoEP first performs an EP-group allgather to collect token-routing information from all ranks, giving every rank the same global routing matrix Ω. The load balancer consumes Ω and generates the replica placement and token routing for the current MoE layer and microbatch. Given the same routing matrix, all ranks independently derive the same plan without a coordinator or result broadcast. The placement and routing decisions remain as device tensors and directly drive the replica-transfer and token-dispatch paths shown in Figure 5. All inputs, intermediate states, and outputs remain in GPU memory, avoiding data-dependent host synchronization between planning and execution. Section 4 details the placement and routing algorithm and its parallel GPU implementation.

3.2

• Main-parameter buffer, with a size of EP · |W | bytes, makes the parameters of main instances owned by the current rank accessible to ranks hosting their replicas. • Replica buffer, with a size of Nslot · |W | bytes, stores the parameters of replicas instantiated on the current rank and is reused to stage and send their gradients during backward propagation.

Replica Buffer Management

• Gradient-accumulation buffer, with a size of E · |W | bytes, receives replica gradients for the current rank’s main instances. The buffer is partitioned by source rank into P regions, each containing EP · |W | bytes.

TopoEP keeps one main instance of each logical expert e on rank main(e), where its model state, optimizer state, and checkpoint ownership remain throughout training. Taking unsharded BF16 mixed-precision training with Adam as an example, a BF16 parameter value, a BF16 gradient, an FP32 master copy, and two FP32 moment estimates together occupy 16 bytes per parameter. A temporary replica receives only the BF16 parameters needed for expert computation and returns its gradients to the main rank for the optimizer update. The optimizer state and checkpoint ownership therefore never migrate with replica placement. Each rank preallocates a fixed number of replica slots and reuses them across layers and microbatches. Applying a new placement only updates the expert-to-slot mapping and copies the selected BF16 parameters. During backward propagation, the same storage is repurposed for replica gradients after the parameters have been consumed. An occupied slot therefore requires only 2 bytes per parameter at any time, one eighth of the complete Adam training state.

3.3

Fine-Grained Pipelined Execution

A standard expert-parallel MoE layer communicates tokens during dispatch and combine. Dynamic replication additionally incurs GPU solver execution, expert-parameter transfers, and replica-gradient transfers. Executing these operations sequentially would place their aggregate latency on the MoE critical path. TopoEP instead partitions the routed tokens into two chunks and schedules communication and expert FFN computation on separate CUDA streams. This pipeline overlaps communication for one chunk with computation for the other, reducing the exposed component of Ttransfer . 3.3.1

Device-initiated replica transfers. At runtime, the generated placement identifies the main

Forward Pass

After the global token-routing all-gather and GPU planning, each rank transfers the parameters required by its expert-to5

TopoEP forward Gate+Plan Comp. Stream

Dispatch+Expert GEMM

the required BF16 parameters from the main rank’s mainparameter buffer into the local replica buffer. The communication stream issues this reload before the parameters are consumed and overlaps it with available computation. The backward pass of each expert GEMM comprises a data-gradient (Dgrad) GEMM and a weight-gradient (Wgrad) GEMM. Dgrad reads the corresponding weight matrix to propagate gradients toward the expert input, whereas Wgrad depends only on the saved activations and output gradients. TopoEP schedules each chunk’s Dgrad before its Wgrad, allowing the resulting token gradients to enter the backward communication path without waiting for weight-gradient computation. The two chunks follow the schedule in Figure 8. While the compute stream executes the expert backward kernels for c1 , the communication stream dispatches the token gradients of c2 . The combine operation for c1 then overlaps the expert backward computation of c2 . The prefetched expert parameters are temporarily staged outside the replica buffer, so the freed buffer is reused to accumulate the Wgrad contributions from c1 and c2 before transferring the result to the main rank’s gradient-accumulation buffer. An EP-group-wide fence ensures that all remote transfers complete before the main rank sums these regions and returns the resulting gradients to autograd. The returned tensors do not alias the accumulation buffer, allowing it to be cleared and reused after the preceding stream operations complete. With a sufficient overlap window, expert-parameter reload and replica-gradient transfer can be hidden by concurrent computation. The exposed incremental overhead then lies primarily in the forward pass, where the GPU solver generates the plan and transfers the selected expert parameters into replica slots. These costs correspond to Tplan and the forward component of Ttransfer , and determine whether the reduction in expert FFN computation time satisfies the end-to-end performance criterion in Eq. 2.

Combine+Next layer

………………………………………………………………… ----------------------- ………………………………………………………………………………………………………………………………-------------------……………

Gate

Plan

idle (wait W + disp c1)

FC1 c1

FC2 c1

idle Res. FC2 c2 (wait comb c2) add

FC1 c2

………………………………………………………………… ----------------------- …………………………………………………………………………………………………………………………………------------------…………

Comm. Stream

…………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………… …………

Pull We

Disp c1

Disp c2

Comb c2

Comb c1

……………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………… ………

time Additional overhead upon standard MoE forward

data-arrival dependency

Overlapped / hidden by computation

output-ready dependency

Figure 7: Two-chunk forward pipeline overlapping expertparameter transfers and token communication with expert FFN computation. TopoEP backward prev compute

Expert bwd · chunk c1

Expert bwd · chunk c2

Gate + Attention bwd

------…………………………………………………………………………………………… Gate idle Attn bwd bwd ………………………………………………………………………………………………………………………………………………………………………………………………………………………-----……………………………………………………………………………………………… ………………………………………………………………………………………………………………………………………………………………………………………………………………………

Comp. Attn bwd Stream (prev layer)

Dgrad c1

Wgrad c1

Dgrad c2

Wgrad c2

prefetched …………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………… …………

Comm. Re-pull We Stream

Disp c1

Disp c2

Comb c1

Comb c2

Push We

……………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………………… ………

time Additional overhead upon standard MoE backward

data-arrival dependency

Overlapped / hidden by computation

output-ready dependency

Figure 8: Two-chunk backward pipeline overlapping expertparameter re-pulls, replica-gradient pushes, and token communication with expert FFN backward computation. slot mapping into the corresponding replica slots. The routed tokens are then partitioned into two disjoint chunks, c1 and c2 . As shown in Figure 7, the communication stream dispatches c1 before its expert FFN computation begins. While the compute stream processes c1 , the communication stream dispatches c2 . It subsequently combines the outputs of c1 while the compute stream processes c2 . This schedule hides the dispatch of c2 and the combine of c1 behind expert FFN computation. The dispatch of c1 and the combine of c2 remain as the pipeline fill and drain boundaries. Compared with the original expert-parallel forward path, the added critical-path latency arises primarily from the GPU solver kernels and expert-parameter transfers. These costs correspond to Tplan and the forward component of Ttransfer , respectively. 3.3.2

4 4.1

Topology-Aware Load-Balancing Algorithm System Model

We consider an expert-parallel MoE training cluster with M NVLink domains. Let D = {0, . . . , M − 1}, R = {0, . . . , R − 1}, and E = {0, . . . , E − 1} denote the sets of domains, EP ranks, and logical experts, respectively. The function dom(r) ∈ D maps rank r to its NVLink domain. We represent the network by a symmetric communication-time matrix C = [cs,r ], where cr,r = 0 and cs,r denotes the effective time to transfer one byte from source rank s to destination rank r. The coefficients can be obtained from profiled effective bandwidth, and are therefore lower for intra-domain NVLink paths than for inter-domain RDMA paths. Each logical expert e has a fixed main instance on rank main(e) ∈ R . This instance holds the expert’s parameters

Backward Pass

Replica slots are reused by subsequent MoE layers during forward propagation, so a layer’s temporary expert parameters are not retained until its backward pass. TopoEP caches the per-layer replica placement and slot mapping, then reloads 6

Replica capacity (C3). Each rank may host at most Nslot non-main expert instances in its preallocated replica slots.

throughout training. Load balancing may create temporary replicas on other ranks by copying these parameters. Let ∥We ∥ denote the parameter size of expert e in bytes and Stok the size in bytes of one routed activation or activation gradient. Each rank reserves Nslot slots for temporary replicas in addition to its fixed main instances. For each microbatch, Top-k gating produces the routing matrix Ω = [ωs,e ] ∈ ZR×E (3) ≥0 ,

∑

ωs,e .

(4)

s∈R dom(s)=d

Since all experts in a layer share the same FFN architecture and have similar per-token compute costs, we estimate each rank’s compute time by multiplying its assigned token count by the profiled average per-assignment time tFFN .

4.2

∀r ∈ R .

Cross-domain replica criterion (C4). Under the two-chunk execution model in Section 3.3, approximately half of the 4Dd,e Stok bytes of forward and backward remote-token traffic remains exposed, while backward replica transfers can be hidden when the overlap window is sufficient. Assuming a common per-byte RDMA cost, a replica of expert e may be placed in a domain different from its main instance only if the forward parameter transfer is smaller than the exposed remote-token traffic:

where ωs,e is the number of token–expert assignments originating from source rank s and selecting expert e. The demand for expert e generated within domain d is Dd,e =

xe,r ≤ Nslot ,

∑

e∈E main(e)̸=r

xe,r = 1,

dom(r) ̸= dom(main(e))

=⇒

∥We ∥ < 2Ddom(r),e Stok .

Fixed main placement (C5). Every expert retains its main instance throughout training. xe,main(e) = 1,

Load-Balancing Formulation

∀e ∈ E .

Given Ω, the optimization jointly determines temporary replica placement and the distribution of each source–expert demand among physical instances of the same logical expert. It preserves every token–expert assignment while minimizing the estimated exposed time of expert computation, token routing, and expert-parameter transfer.

Objective. The exposed expert FFN computation time is

Decision variables. We introduce two decision variables. The binary variable xe,r indicates whether rank r hosts an instance of expert e, including its main instance. The nonnegative integer qs,e,r gives the number of (s, e) assignments executed on destination rank r.

Ttoken = 2Stok ∑ cs,r qs,e,r .

Tcomp = tFFN Lmax .

Under the overlap assumptions in (C4), the exposed tokenrouting and forward expert-parameter transfer times are

= ∑ ∑ qs,e,r ,

Treplica = ∑

e∈E

cmain(e),r ∥We ∥xe,r .

(7)

min Tcomp + Ttoken + Treplica

∀r ∈ R ,

x,q

s.t.

Lmax = max Lr .

(C1)–(C5),

(8)

xe,r ∈ {0, 1},

r∈R

qs,e,r ∈ Z≥0 ,

Constraints. A feasible replica-placement and tokenrerouting plan satisfies the following conditions. Token conservation (C1). Every token–expert assignment is mapped to exactly one physical instance.

4.3

∀e ∈ E , s, r ∈ R .

Two-Stage GPU-Native Solver

Solving Eq. 8 exactly for every layer and microbatch would add excessive latency to the training critical path. TopoEP therefore constructs a feasible placement-and-routing plan with a topology-aware two-stage GPU solver. Inter-node placement first creates replicas when the reduction in crossdomain token traffic justifies the parameter transfer, and intranode refinement then reduces the remaining rank-load imbalance within each NVLink domain.

∀s ∈ R , e ∈ E .

r∈R

Instance reachability (C2). Tokens may be assigned only to ranks that host the corresponding expert. qs,e,r ≤ ωs,e xe,r ,

∑

r∈R r̸=main(e)

The optimization minimizes their sum as an idealized estimate of exposed execution time:

s∈R e∈E

∑ qs,e,r = ωs,e ,

(6)

s,r∈R e∈E

Load metric. The assigned load of rank r is the total number of token–expert assignments executed on that rank. We define Lmax as the maximum assigned load across all ranks. Lr

(5)

∀s, r ∈ R , e ∈ E . 7

NVL 0

GPU 0

E0

GPU 1

E1

E2

Input: Expert Demand (Ω) Expert

RDMA

E3

GPU 0

Tokens destination 1

GPU0 · Local

E1

1

GPU0 · Local

E4

8

GPU2 · RDMA

E0

Load = 14

Expert Tokens destination 2

GPU0 · NVLink

E1

2

GPU0 · NVLink

E2

3

GPU1 · Local

E3

3

GPU1 · Local

local token

E7

NVL 0

NVL 1 GPU 1

E4’

E6

E2

GPU 0

GPU 2 E4

E3

E0

GPU 1 E4’

E1

E2

E3

E4’’

……

GPU1 routes E0

E1

E5

GPU 3

Stage 2: Intra-node Refinement

NVL 0

E0

E4

Stage 1: Inter-node Placement

GPU0 routes

NVL 1

GPU 2

Load = 6

Load = 10

Load = 10

Incoming tokens

……

token via NVLink

E4

Main expert

E4’/ E4’’

4 E4 tokens via NVLink

Replica

Free slot

NVLink

Figure 9: Example of two-stage expert replication and token rerouting. Inter-node placement removes cross-domain token traffic, after which intra-node refinement balances rank loads through same-domain rerouting over NVLink. Figure 9 illustrates the two stages using two NVLink domains. In the input routing, eight assignments from GPU 0 select expert E4 , whose main instance resides on GPU 2. These assignments therefore cross the inter-node RDMA fabric during both dispatch and combine. Inter-node placement creates replica E4′ on GPU 0 when one parameter transfer costs less than the exposed bidirectional token traffic. The eight assignments can then execute locally, eliminating their RDMA token traffic. This placement, however, leaves GPU 0 and GPU 1 with loads of 14 and 6. Intra-node refinement creates replica E4′′ on GPU 1 and moves four E4 assignments from GPU 0 to GPU 1 over NVLink, balancing both rank loads at 10 without introducing new cross-domain token routes.

cal token–expert assignment while preferring same-domain destinations. Intra-node refinement. Starting from the inter-node placement, each NVLink domain independently reduces its residual rank-load imbalance. In every iteration, the solver selects the busiest rank rb and evaluates same-domain target ranks rt for the experts contributing to its load. A target is feasible if it already hosts the expert or has a free replica slot. The candidate transfer amount is    Lrb − Lrt δ = min Ue,rb , . (10) 2 The half-gap bound reduces the difference between the source and target loads without reversing their order. Candidates are ranked by transfer amount, target load, and communication cost. The selected update creates a replica when needed and modifies x, q, U, and L. Because the target lies in the same domain as the overloaded rank, refinement improves load balance without placing the expert in a new domain or adding cross-domain token routes.

Inter-node placement. For every domain d, the solver aggregates the routing matrix by source domain and computes the replication benefit bd,e = 2Dd,e Stok − ∥We ∥.

(9)

For an expert whose main instance lies outside domain d, a positive bd,e means that replacing the exposed remote-token traffic with one parameter transfer is beneficial. Each domain considers its positive-benefit candidates in descending order and places replicas on the least occupied ranks with free slots. Processing domains independently allows these decisions to run concurrently while enforcing the per-rank slot limit. Given the resulting placement, the solver constructs the physical token routing. For each source–expert pair (s, e), it selects instances in the source domain whenever they are available and otherwise uses all deployed instances of e. It divides ωs,e approximately evenly among the selected instances. It then computes the per-instance loads Ue,r = ∑s qs,e,r and updates the rank loads Lr . This routing preserves every logi-

Constraint preservation. Inter-node placement starts from the fixed main instances, creates replicas only on ranks with free slots, and admits a cross-domain expert–domain pair only when bd,e > 0. It therefore satisfies replica capacity (C3), the cross-domain replica criterion (C4), and fixed main placement (C5). Routing partitions every ωs,e completely among deployed instances, satisfying token conservation (C1) and instance reachability (C2). Intra-node refinement only transfers existing assignments to an existing instance or a new replica in a free slot, and adds replicas only within a domain that already hosts the expert. Each refinement step therefore preserves (C1)–(C5). 8

GPU parallelism and determinism. Inter-node candidate evaluation and source–expert routing run concurrently across CUDA thread blocks. During intra-node refinement, one block handles each NVLink domain, caches its placement and load state in shared memory, and uses warp- and block-level reductions to select a candidate. The block commits one update per iteration to avoid conflicting modifications. Deterministic candidate ordering allows every rank to generate the same x and q from the shared routing matrix, eliminating the need for a coordinator or plan broadcast.

5

Table 1: MoE model configurations used in the evaluation.

Qwen3-30B-A3B [31] GLM-4.5-Air [36] DeepSeek-V2 [3]

GPU-native solver. Inter-node placement, routing update, and intra-node refinement are implemented as CUDA kernels. Independent domains and source–expert tasks execute concurrently, with warp- and block-level reductions used for candidate selection.

DeepEP dispatch and combine. We integrate DeepEP V2 [38] (commit af9a040) into our communication manager to replace Megatron-LM’s NCCL-based all-to-all path. The manager consumes device-resident routing information and derives receive counts on the GPU, avoiding CPU synchronization during dispatch and combine. Fine-grained two-chunk pipeline. We use a custom autograd function to explicitly schedule the expert computation and communication of each MoE layer across separate CUDA streams. The pipeline overlaps token communication and replica transfers with expert computation, reducing their exposed critical-path overhead.

6.2

Experimental Setup

8 8 6

2,048 4,096 5,120

768 1,408 1,536

1/32 2/16 2/16

Models and Workloads. We evaluate TopoEP with Qwen330B-A3B [31], GLM-4.5-Air [36], and DeepSeek-V2 [3]. As summarized in Table 1, these models span 128–160 routed experts, Top-k values of 6 and 8, and substantially different hidden and expert-intermediate dimensions. Qwen3 uses PP 1/EP 32, while GLM-4.5-Air and DeepSeek-V2 use PP 2/EP 16. To fit the models within the available GPU memory, we reduce the number of Transformer layers while retaining their original MoE-layer dimensions and routing configurations. This setup keeps the evaluation focused on MoE-layer performance. We construct a training corpus from six public sources: English web text from FineWeb [23]; Chinese web text from FineWeb2 [24]; diverse pretraining text from Dolma [29]; source code in Python, C++, Java, and Rust from StarCoderData [15]; scientific literature from peS2o [30]; and mathematical reasoning prompts from DAPO-Math-17K [35]. We deduplicate and truncate documents to at most 4,096 tokens before assembling them into training sequences.

Device-initiated replica transfers. Main-parameter, replicaslot, and gradient-accumulation buffers are allocated and registered in NCCL symmetric-memory windows during initialization. TMA handles intra-node transfers over NVLink, while NCCL GIN performs device-initiated get/put operations for inter-node RDMA.

6.1

128 128 160

Baselines. We compare TopoEP with Megatron-LM [27] and three representative MoE load-balancing methods: DeepSeekEPLB [5], the official expert-replication and placement heuristic; FasterMoE [10], which uses dynamic expert shadowing; and FlexMoE [21], which expands, shrinks, and migrates virtual experts. We adapt each method’s core load-balancing algorithm to the same Megatron-LM codebase at commit 0ff7226. Within each comparison, all dynamic methods use the same additional replica budget, and all runs use the same routing configuration.

We integrate TopoEP into Megatron-LM [27]. The implementation comprises approximately 10,000 lines of Python and CUDA/C++.

Performance Evaluation

5 5 4

supports 32-way expert parallelism and includes both intranode and inter-node communication paths.

Implementation

6

MoE Interm. PP / EP Experts Top-k Hidden layers size

Model

Load-Balancing Effectiveness

Figure 10 compares the maximum-to-mean rank-load ratio as the additional replica budget increases. Without balancing, the ratio is 8.39. With one slot, the ratio for TopoEP remains 3.79, whereas FlexMoE reaches 1.45. With two slots, TopoEP drops to 1.30 and outperforms FlexMoE, FasterMoE, and DeepSeekEPLB by 6.8%, 48.2%, and 55.8%, respectively. With three or four slots, TopoEP maintains a rank-load ratio of 1.30, while FasterMoE and DeepSeek-EPLB remain at 1.82 and 2.88 even with four slots. Thus, TopoEP needs only two additional slots per rank to achieve the best load balance among the evaluated methods; allocating more slots provides almost no further

Testbed. Our experiments run on a four-node cluster with 32 NVIDIA H800 GPUs. Each node contains eight GPUs interconnected through NVLink and NVSwitch, providing up to 400 GB/s of aggregate bidirectional bandwidth. A railoptimized InfiniBand fabric provides inter-node communication, with each GPU connected to a dual-port NVIDIA ConnectX-7 NIC through two 200 Gb/s links. The testbed 9

8.39

Ideal No Balance

Rank-load imbalance (max / mean)

8

DeepSeek-EPLB FasterMoE

6.4

FlexMoE TopoEP

In this section, we compare the end-to-end training performance of TopoEP against multiple baselines and examine how routing skew affects training throughput.

7

6

5

Sensitivity to routing skew. Natural routing uses the router logits generated by the model from the input tokens. As discussed in Section 2.2, expert-load imbalance commonly arises in MoE training and gives load balancing greater opportunity to reduce straggler delays. Because expert specialization typically emerges over long pretraining runs, our reduceddepth models and limited training duration may not exhibit the full range of routing imbalance that can arise at full scale. We therefore construct synthetic routing workloads with controlled skew to test whether TopoEP mitigates the resulting throughput degradation across a broader range of imbalance. Synthetic routing uses random logits with a fixed expert-specific bias. A router_skew value of 0 produces uniform expert selection in expectation, while increasingly negative values strengthen the bias and concentrate more tokens on a subset of experts. We use these synthetic settings only for performance measurements. Under natural routing, Figure 12 shows that TopoEP improves throughput over Megatron-LM by 9.5%, 11.4%, and 6.2% on Qwen, GLM, and DeepSeek-V2, respectively. Under uniform synthetic routing (router_skew= 0), it instead trails Megatron-LM by 8.5%– 16.7%. Thus, load balancing does not improve every routing workload. When expert demand is already uniform, the limited reduction in straggling does not offset planning and replica-management costs. As synthetic routing becomes more imbalanced from router_skew= 0 to −4, Megatron-LM throughput falls by 27.8%, 25.8%, and 24.7% across the three models. In contrast, TopoEP throughput varies by less than 2.3% across all synthetic settings. At router_skew= −4, TopoEP outperforms Megatron-LM by 26.6%, 20.9%, and 10.0%. These results show that its throughput benefit grows as routing imbalance increases. More importantly, TopoEP maintains nearly constant throughput across the full skew range. This robustness shows that substantial changes in routing and expert-load imbalance need not translate into large performance fluctuations, enabling stable training performance as routing behavior evolves.

4

3

2

1

0

1

No balance

2

3

4

Replica slots per rank (Nslot)

Figure 10: Rank-load imbalance versus additional replica slots per rank for a 32-rank, 640-expert Top-8 workload. Lower is better, and the dashed line denotes ideal balance. Kernel solve latency (us)

End-to-End Performance

510 us

500 402 us

455 us

415 us

400

414 us

360 us

429 us

337 us 368 us

300

431 us

357 us

369 us

200 131 us

100

8

16

32

64

8 12

64

8 12

6 25

4 2 0 8 4 38 51 64 76 102

EP ranks

Expert nums

(a) Rank scaling (E=640)

(b) Expert scaling (R=32)

Figure 11: Mean solver-kernel latency with varying logical EP sizes and expert counts, measured over 200 executions after 20 warm-up iterations.

benefit. This small replica budget also reduces the additional GPU memory overhead introduced by expert replication. We use this setting in the remaining experiments.

6.3

Solver Scalability

Figure 11(a) shows that, with the number of experts fixed at E = 640, kernel latency increases from 131 µs at EP size 8 to 510 µs at EP size 128. At EP size 8, all ranks fit within a single NVLink domain, so the solver does not need to evaluate inter-node placement, contributing to the lower kernel latency. Figure 11(b) varies the expert count by 16× at a fixed EP size of 32, while latency remains between 337 and 455 µs. Increasing the EP size by 16× raises latency by only 3.9×, while the same increase in expert count changes the endpoint latency by only 1.27×. Thus, the problem grows by an order of magnitude without a proportional latency increase. Even the largest tested configuration is solved in just 0.510 ms, preserving microsecond-scale planning at every layer and microbatch.

Comparison with baselines. To ensure a fair comparison of load-balancing policies, all evaluated EPLB methods use the same optimized plan-execution backend described in Section 5, including TMA / GIN replica transfers, DeepEP token dispatch and combine, and the two-chunk execution pipeline. They differ only in the planner that generates the placement and routing decisions; Megatron-LM serves separately as the no-EPLB baseline. Figure 13(a) shows that TopoEP achieves the highest throughput among all evaluated methods on all three models. It exceeds the highest baseline throughput by 4.4% on Qwen, 10

Average throughput (token/s)

Megatron-LM

TopoEP

120k

450k

120k 110k

400k

110k 100k

1.27 ×

1.10 ×

1.21 × 100k

350k

90k 0

Natural

−0.5 −1

−4

−2

0

Natural

−0.5 −1

−2

−4

0

Natural

−0.5 −1

−2

Routing policy

Routing policy

Routing policy

(a) Qwen3-30B-A3B

(b) GLM-4.5-Air

(c) DeepSeek-V2

−4

Figure 12: Average end-to-end training throughput under natural and synthetic routing. The synthetic settings vary router_skew from 0 to −4. Stars mark interpolated throughput crossovers, and double-headed arrows show the TopoEP/Megatron-LM throughput ratio at router_skew= −4.

80

80

60

60

40

40

20

20

0

0

103.5

97.3

96.8

94.1

93.2

120 100

Megatron-LM

TopoEP

Language-model loss

300

92.5

100

88.8

88.3

120

FlexMoE

99.3

DeepSeek-EPLB

106.7

FasterMoE

417.5

399.8

347.8

335.5

400

329.7

Average throughput (k tokens/s)

Megatron-LM

200

0

Qwen3-30B-A3B

GLM-4.5-Air

6

1

1k

2k

3k 7k

8k

9k

10k

Training step

Figure 14: Qwen3-30B-A3B training loss during 10,000-step runs under natural routing.

2.52 1.03

1.05

1

8

TopoEP

2.69 1.98

1.41

2

10

DeepSeek-V2

DeepSeek-EPLB

4.30

FlexMoE

2.16

2.48

3.24

4 3

GLM-4.5-Air

FasterMoE

2.97

5

4.56

Megatron-LM

3.88

Qwen3-30B-A3B

3.07

0

4.09

Average rank-load imbalance (max / mean)

100

TopoEP

trajectories for TopoEP and Megatron-LM under the same model initialization, data order, and learning-rate schedule. Across all 10,000 steps, the mean and maximum absolute relative differences are 0.088% and 0.407%, respectively. These results show no observable degradation in loss convergence.

DeepSeek-V2

Model

Figure 13: Comparison with baselines at router_skew= −4. (a) Average end-to-end throughput. (b) Average rank-load imbalance, defined as the maximum per-rank token-expert load divided by the average across ranks, where 1 denotes ideal balance.

6.5

Latency Breakdown

Figure 15 shows that TopoEP reduces expert-computation time on straggler ranks by 49.2%–64.0% across the forward and backward passes of Qwen and GLM, directly demonstrating its effectiveness in mitigating computation stragglers. Token dispatch/combine time also decreases by 34.9%–65.3%. TopoEP introduces 1.04–1.58 ms of replica-management time per layer, which is smaller than the reduction in each of the other two components in every panel. Because these components may overlap across CUDA streams, their measured times cannot be added directly. Overall, the breakdown shows that TopoEP reduces both expert-computation and token dispatch/combine time under imbalanced routing.

14.5% on GLM, and 4.2% on DeepSeek-V2. Figure 13(b) shows that TopoEP also achieves the lowest average rankload imbalance, reducing the ratios to 1.41, 1.05, and 1.03, respectively. These values are 34.7%, 47.1%, and 59.2% lower than the lowest baseline ratios. With the plan-execution backend held constant across EPLB methods, these results show that TopoEP achieves better rank-load balance while also delivering higher endto-end throughput when planner and execution overheads are included. Training stability. Figure 14 shows closely overlapping loss 11

5

7.64

10

0 Megatron-LM Baseline

TopoEP

7.60 16.02

(c) GLM-4.5-Air Forward

7.42

0 Megatron-LM Baseline

TopoEP

39.5

50.0

37.5

1.9

1.1

6.8

18.5

38.0

50.0

43.8

54.9

75.0

40 20 0

Qwen3-30B-A3B

68.3

54.7

63.8

34.8

60

GLM-4.5-Air

DeepSeek-V2

(b) Assignments served by replicas

Figure 16: Physical execution of token–expert assignments. (a) Fraction whose source and execution ranks are on different nodes. (b) Fraction executed by replicas.

Topology-Aware Token Routing and DeepSeek-EPLB [5] places redundant experts using a heuristic based on estimated historical loads. These approaches can improve expert utilization, but host-side planning and reconfiguration add CPU–GPU coordination and state-management overhead, making fine-grained adaptation hard to run on the critical path of each microbatch.

We examine whether the plans generated by the algorithm in Section 4 are topology-aware in practice. All methods receive the same logical Top-k assignments, with a common replica budget for the EPLB methods. They differ in the placement of expert instances and the physical ranks selected to execute each assignment. For each token–expert assignment, we record whether its execution rank is on a different node from its source rank and whether it is served by a replica. Figure 16(a)–(b) shows that TopoEP reduces the inter-node assignment fraction from 50.0%–75.0% with Megatron-LM to 0.91%–1.88%, while replicas execute 63.8%–79.8% of assignments. Across all models, it achieves the lowest internode fraction and the highest replica-served fraction. Together, these results indicate that TopoEP redistributes most assignments to replicas while keeping their execution within the source nodes. The comparison with DeepSeek-EPLB shows that high replica usage alone does not ensure communication locality. On GLM-4.5-Air, DeepSeek-EPLB and TopoEP serve 58.6% and 63.8% of assignments through replicas, respectively, yet their inter-node fractions are 18.5% and 1.1%. This contrast highlights the importance of coordinating replica placement with token allocation, consistent with TopoEP’s two-stage strategy of placing expert instances near token sources and refining rank loads within each node.

7

80

(d) GLM-4.5-Air Backward

Figure 15: Per-layer MoE latency breakdown, with expertcomputation time measured on straggler ranks.

6.6

DeepSeek-V2

23.4

1.58

8.00 2.99

GLM-4.5-Air

(a) Inter-node assignments

11.67

1.32

TopoEP

0.0

20

10

Qwen3-30B-A3B

58.6

Time (ms)

20 15.40

0

(b) Qwen3-30B-A3B Backward

30

25

15

TopoEP

20

17.5

Megatron-LM Baseline

(a) Qwen3-30B-A3B Forward

3.51

11.7

TopoEP

0

40

0.0

Megatron-LM Baseline

1.37

4.47 6.91

DeepSeek-EPLB

79.8

3.81

0

5

69.4

4.55

60

52.2

1.04

5

1.07

FlexMoE

0.9

10

39.9

10.28

13.11

80

65.6

15

FasterMoE

12.1

10

Megatron-LM

0.0

15

Replica Management

Fraction of assignments (%)

Token Dispatch/Combine

Fraction of assignments (%)

Time (ms)

Expert Compute

GPU-native load balancing. More recent systems move load-balancing decisions onto the GPU. UltraEP [34] and MoonEP [1] use lightweight planners for per-layer, permicrobatch balancing and avoid the host–GPU round trip. However, both are designed primarily for a single scale-up domain. Their planners neither model the heterogeneous communication costs of NVLink and RDMA nor optimize crossdomain replica placement, so they are not directly applicable to the inter-node EP setting considered here. TopoEP targets this setting with a planner that models the hierarchical scaleup and scale-out topology when deciding replica placement and token rerouting.

8

Conclusion

This paper presents TopoEP, a GPU-native, topology-aware load-balancing system for large-scale MoE training. It uses a deterministic two-stage GPU solver to balance rank loads while accounting for the heterogeneous communication costs of scale-up and scale-out fabrics, and applies each resulting plan without data-dependent host synchronization. When integrated with Megatron-LM on a 32-GPU NVIDIA H800 cluster, it improves end-to-end training throughput by 6.2%– 11.4% across three representative MoE models without observable degradation in training-loss convergence.

Related Work

Host-side load balancing. System-level EPLB designs commonly rely on host-side planning. FasterMoE [10] replicates overloaded experts via shadow experts, FlexMoE [21] dynamically expands, shrinks, and migrates expert replicas, 12

References

[11] Changho Hwang, Wei Cui, Yifan Xiong, Ziyue Yang, Ze Liu, Han Hu, Zilong Wang, Rafael Salas, Jithin Jose, Prabhat Ram, HoYuen Chau, Peng Cheng, Fan Yang, Mao Yang, and Yongqiang Xiong. Tutel: Adaptive mixture-of-experts at scale. In Proceedings of Machine Learning and Systems, volume 5, pages 269–287. MLSys, 2023.

[1] Yutian Chen, Cong Li, Yucheng Wang, and Ming Wei. MoonEP: A perfectly balanced expert parallelism library via dynamic redundant experts. https://github.com/ MoonshotAI/MoonEP, 2026. [2] Damai Dai, Chengqi Deng, Chenggang Zhao, R. X. Xu, Huazuo Gao, Deli Chen, Jiashi Li, Wangding Zeng, Xingkai Yu, Y. Wu, Zhenda Xie, Y. K. Li, Panpan Huang, Fuli Luo, Chong Ruan, Zhifang Sui, and Wenfeng Liang. DeepSeekMoE: Towards ultimate expert specialization in mixture-of-experts language models. In Proceedings of the 62nd Annual Meeting of the Association for Computational Linguistics (Volume 1: Long Papers), pages 1280–1297. Association for Computational Linguistics, 2024.

[12] Yechan Kim, Hwijoon Lim, and Dongsu Han. Scaling beyond the GPU memory limit for large mixture-ofexperts model training. In Ruslan Salakhutdinov, Zico Kolter, Katherine Heller, Adrian Weller, Nuria Oliver, Jonathan Scarlett, and Felix Berkenkamp, editors, Proceedings of the 41st International Conference on Machine Learning, volume 235 of Proceedings of Machine Learning Research, pages 24342–24353. PMLR, 21–27 Jul 2024. [13] 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 9th International Conference on Learning Representations (ICLR), 2021.

[3] DeepSeek-AI. Deepseek-v2: A strong, economical, and efficient mixture-of-experts language model, 2024. [4] DeepSeek-AI. DeepSeek-V3 technical report, 2024. [5] DeepSeek-AI. EPLB: Expert parallelism load balancer. https://github.com/deepseek-ai/EPLB, 2025. [6] 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.

[14] Jiamin Li, Yimin Jiang, Yibo Zhu, Cong Wang, and Hong Xu. Accelerating distributed MoE training and inference with Lina. In 2023 USENIX Annual Technical Conference (USENIX ATC 23), pages 945–959, Boston, MA, July 2023. USENIX Association.

[7] Trevor Gale, Deepak Narayanan, Cliff Young, and Matei Zaharia. MegaBlocks: Efficient sparse training with mixture-of-experts. In Proceedings of Machine Learning and Systems, volume 5, pages 288–304. MLSys, 2023.

[15] Raymond Li, Loubna Ben Allal, Yangtian Zi, Niklas Muennighoff, Denis Kocetkov, Chenghao Mou, Marc Marone, Christopher Akiki, et al. StarCoder: May the source be with you! Transactions on Machine Learning Research, 2023. [16] Juncai Liu, Jessie Hui Wang, and Yimin Jiang. Janus: A unified distributed training framework for sparse mixture-of-experts models. In Proceedings of the ACM SIGCOMM 2023 Conference, pages 486–498. ACM, 2023.

[8] Wentao Guo, Mayank Mishra, Xinle Cheng, Ion Stoica, and Tri Dao. Sonicmoe: Accelerating moe with io and tile-aware optimizations. In C. Vondrick, B. Hariharan, C. Raffel, L. Pinto, D. Yang, and A. Faust, editors, International Conference on Learning Representations, volume 2026, pages 67814–67839, 2026.

[17] Rui Liu, Young Jin Kim, Alexandre Muzio, and Hany Hassan. Gating dropout: Communication-efficient regularization for sparsely activated transformers. In Kamalika Chaudhuri, Stefanie Jegelka, Le Song, Csaba Szepesvari, Gang Niu, and Sivan Sabato, editors, Proceedings of the 39th International Conference on Machine Learning, volume 162 of Proceedings of Machine Learning Research, pages 13782–13792. PMLR, 17–23 Jul 2022.

[9] Khaled Hamidouche, John Bachan, Pak Markthub, PeterJan Gootzen, Elena Agostini, Sylvain Jeaugey, Aamir Shafi, Georgios Theodorakis, and Manjunath Gorentla Venkata. Gpu-initiated networking for nccl. arXiv preprint arXiv:2511.15076, 2025. [10] Jiaao He, Jidong Zhai, Tiago Antunes, Haojie Wang, Fuwen Luo, Shangfeng Shi, and Qin Li. FasterMoE: Modeling and optimizing training of large-scale dynamic pre-trained models. In Proceedings of the 27th ACM SIGPLAN Symposium on Principles and Practice of Parallel Programming, pages 120–134. ACM, 2022.

[18] Zixuan Ma, Jiaao He, Jiezhong Qiu, Huanqi Cao, Yuanwei Wang, Zhenbo Sun, Liyan Zheng, Haojie Wang, Shizhi Tang, Tianyu Zheng, Junyang Lin, Guanyu Feng, Zeqiang Huang, Jie Gao, Aohan Zeng, Jianwei Zhang, 13

Runxin Zhong, Tianhui Shi, Sha Liu, Weimin Zheng, Jie Tang, Hongxia Yang, Xin Liu, Jidong Zhai, and Wenguang Chen. Bagualu: targeting brain scale pretrained models with over 37 million cores. In Proceedings of the 27th ACM SIGPLAN Symposium on Principles and Practice of Parallel Programming, PPoPP ’22, page 192–204, New York, NY, USA, 2022. Association for Computing Machinery.

Dean. Outrageously large neural networks: The sparselygated mixture-of-experts layer. In 5th International Conference on Learning Representations (ICLR), 2017. [27] Mohammad Shoeybi, Mostofa Patwary, Raul Puri, Patrick LeGresley, Jared Casper, and Bryan Catanzaro. Megatron-lm: Training multi-billion parameter language models using model parallelism, 2020. [28] Athinagoras Skiadopoulos, Mark Zhao, Swapnil Gandhi, Thomas Norrie, Shrijeet Mukherjee, and Christos Kozyrakis. SYMI: Efficient Mixture-of-Experts training via model and optimizer state decoupling. In 23rd USENIX Symposium on Networked Systems Design and Implementation (NSDI 26), pages 75–92, Renton, WA, May 2026. USENIX Association.

[19] Xuan-Phi Nguyen, Shrey Pandit, Austin Xu, Caiming Xiong, and Shafiq Joty. Least-loaded expert parallelism: Load balancing an imbalanced mixture-of-experts. In Forty-third International Conference on Machine Learning, 2026. [20] Xiaonan Nie, Xupeng Miao, Shijie Cao, Lingxiao Ma, Qibin Liu, Jilong Xue, Youshan Miao, Yi Liu, Zhi Yang, and Bin Cui. Evomoe: An evolutional mixture-ofexperts training framework via dense-to-sparse gate, 2022.

[29] Luca Soldaini, Rodney Kinney, Akshita Bhagia, Dustin Schwenk, David Atkinson, Russell Authur, Ben Bogin, Khyathi Chandu, Jennifer Dumas, Yanai Elazar, Valentin Hofmann, Ananya Jha, Sachin Kumar, Li Lucy, Xinxi Lyu, Nathan Lambert, Ian Magnusson, Jacob Morrison, Niklas Muennighoff, Aakanksha Naik, Crystal Nam, Matthew Peters, Abhilasha Ravichander, Kyle Richardson, Zejiang Shen, Emma Strubell, Nishant Subramani, Oyvind Tafjord, Evan Walsh, Luke Zettlemoyer, Noah Smith, Hannaneh Hajishirzi, Iz Beltagy, Dirk Groeneveld, Jesse Dodge, and Kyle Lo. Dolma: an open corpus of three trillion tokens for language model pretraining research. In Lun-Wei Ku, Andre Martins, and Vivek Srikumar, editors, Proceedings of the 62nd Annual Meeting of the Association for Computational Linguistics (Volume 1: Long Papers), pages 15725–15788, Bangkok, Thailand, August 2024. Association for Computational Linguistics.

[21] Xiaonan Nie, Xupeng Miao, Zilong Wang, Zichao Yang, Jilong Xue, Lingxiao Ma, Gang Cao, and Bin Cui. FlexMoE: Scaling large-scale sparse pre-trained model training via dynamic device placement. Proceedings of the ACM on Management of Data, 1(1):1–19, 2023. [22] NVIDIA. NVIDIA Hopper Architecture InDepth. https://developer.nvidia.com/blog/ nvidia-hopper-architecture-in-depth/, 2022. [23] Guilherme Penedo, Hynek Kydlíček, Loubna Ben allal, Anton Lozhkov, Margaret Mitchell, Colin Raffel, Leandro Von Werra, and Thomas Wolf. The fineweb datasets: Decanting the web for the finest text data at scale. In The Thirty-eight Conference on Neural Information Processing Systems Datasets and Benchmarks Track, 2024.

[30] Luca Soldaini and Kyle Lo. peS2o (Pretraining Efficiently on S2ORC) Dataset. Technical report, Allen Institute for AI, 2023. ODC-By, https://github.com/ allenai/pes2o.

[24] Guilherme Penedo, Hynek Kydlíček, Vinko Sabolčec, Bettina Messmer, Negar Foroutan, Amir Hossein Kargaran, Colin Raffel, Martin Jaggi, Leandro von Werra, and Thomas Wolf. FineWeb2: One pipeline to scale them all—adapting pre-training data processing to every language, 2025.

[31] Qwen Team. Qwen3 technical report, 2025. [32] Ashish Vaswani, Noam Shazeer, Niki Parmar, Jakob Uszkoreit, Llion Jones, Aidan N. Gomez, Łukasz Kaiser, and Illia Polosukhin. Attention is all you need. In Advances in Neural Information Processing Systems, volume 30, pages 5998–6008. Curran Associates, Inc., 2017.

[25] Samyam Rajbhandari, Conglong Li, Zhewei Yao, Minjia Zhang, Reza Yazdani Aminabadi, Ammar Ahmad Awan, Jeff Rasley, and Yuxiong He. DeepSpeed-MoE: Advancing mixture-of-experts inference and training to power next-generation AI scale. In Proceedings of the 39th International Conference on Machine Learning, volume 162 of Proceedings of Machine Learning Research, pages 18332–18346. PMLR, 2022.

[33] Lean Wang, Huazuo Gao, Chenggang Zhao, Xu Sun, and Damai Dai. Auxiliary-loss-free load balancing strategy for mixture-of-experts. arXiv preprint arXiv:2408.15664, 2024.

[26] Noam Shazeer, Azalia Mirhoseini, Krzysztof Maziarz, Andy Davis, Quoc V. Le, Geoffrey E. Hinton, and Jeff

[34] Xinming Wei, Chao Jin, Tuo Dai, Yinmin Zhong, Shan Yu, Chengxu Yang, Bingyang Wu, Zili Zhang, Jing Mai, 14

Qianchao Zhu, Zhouyang Li, Yuliang Liu, and Guojie Luo. UltraEP: Unleash MoE training and inference on rack-scale nodes with near-optimal load balancing, 2026. [35] Qiying Yu, Zheng Zhang, Ruofei Zhu, Yufeng Yuan, Xiaochen Zuo, Yu Yue, et al. DAPO: An open-source LLM reinforcement learning system at scale. In D. Belgrave, C. Zhang, H. Lin, R. Pascanu, P. Koniusz, M. Ghassemi, and N. Chen, editors, Advances in Neural Information Processing Systems, volume 38, Main Conference, pages 113222–113244. Curran Associates, Inc., 2025. [36] Aohan Zeng, Xin Lv, Qinkai Zheng, Zhenyu Hou, Bin Chen, et al. GLM-4.5: Agentic, reasoning, and coding (ARC) foundation models, 2025. [37] Shulai Zhang, Ningxin Zheng, Haibin Lin, Ziheng Jiang, Wenlei Bao, Chengquan Jiang, Qi Hou, Weihao Cui, Size Zheng, Li-Wen Chang, Quan Chen, and Xin Liu. Comet: Fine-grained computation-communication overlapping for mixture-of-experts. In M. Zaharia, G. Joshi, and Y. Lin, editors, Proceedings of Machine Learning and Systems, volume 7. MLSys, 2025. [38] 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. [39] Yanqi Zhou, Tao Lei, Hanxiao Liu, Nan Du, Yanping Huang, Vincent Zhao, Andrew M Dai, Quoc V Le, James Laudon, et al. Mixture-of-experts with expert choice routing. Advances in Neural Information Processing Systems, 35:7103–7114, 2022.

15

Record · ID 1108697 · SHA-256 35c003b01905f5ce
Retrieved via Conceptio — every document is proof-bundled with source, license, and retrieval metadata.