作者Zhiyi Yao, Zuning Liang, Yuedong Xu, Jin Zhao, Jessie Hui Wang, Tong Li
With the ever-growing size of deep learning models, GPU memory is prone to being insufficient during training. A prominent approach is ZeRO-Offload, which moves the optimizer states to CPU memory and performs parameter update using CPU. However, the deficiencies of ZeRO-Offload include low GPU utilization, imperfect overlapping of communication and computation, and inflexible offloading. In this paper, we leverage Direct Host Access (DHA) on the GPU that can compute data in CPU memory, forming a novel hybrid on-GPU and DHA. We design and implement MemFerry consisting of an execution scheduler and a shadow model. The scheduler strategically chooses layers of parameters for DHA computation and transmits the remaining parameters to GPU memory simultaneously to shorten forward propagation time, and further loads DHA parameters to GPU memory to reduce backward propagation time. The shadow model presents a unified memory abstraction for the parameter partitions stored separately in GPU and CPU memories. To further reduce GPU memory usage, we present MemFerry along with its dynamic programming algorithm that offloads gradients to CPU memory via DHA. We further extend MemFerry to emerging scale-up domains with ScaleUp-MemFerry, which exploits otherwise underutilized accelerator interconnect bandwidth to assist host-to-accelerator data movement through adaptive multi-path transfer. Our experiments show that \system trains up to $1.68\times$ faster and MemFerry can train $1.52\times$ larger model compared to ZeRO-Offload on a single GPU, and increase training speed by at least $28.1%$ when scaling to data parallelism on 8 GPUs. We further extend the design to a Huawei CloudMatrix384 scale-Up node with up to 8 NPUs, and our ScaleUp-MemFerry reduces the end-to-end iteration time by up to $20.7%$ over DeepSpeed.
展开完整摘要收起摘要↓
With the ever-growing size of deep learning models, GPU memory is prone to being insufficient during training. A prominent approach is ZeRO-Offload, which moves the optimizer states to CPU memory and performs parameter update using CPU. However, the deficiencies of ZeRO-Offload include low GPU utilization, imperfect overlapping of communication and computation, and inflexible offloading. In this paper, we leverage Direct Host Access (DHA) on the GPU that can compute data in CPU memory, forming a novel hybrid on-GPU and DHA. We design and implement MemFerry consisting of an execution scheduler and a shadow model. The scheduler strategically chooses layers of parameters for DHA computation and transmits the remaining parameters to GPU memory simultaneously to shorten forward propagation time, and further loads DHA parameters to GPU memory to reduce backward propagation time. The shadow model presents a unified memory abstraction for the parameter partitions stored separately in GPU and CPU memories. To further reduce GPU memory usage, we present MemFerry along with its dynamic programming algorithm that offloads gradients to CPU memory via DHA. We further extend MemFerry to emerging scale-up domains with ScaleUp-MemFerry, which exploits otherwise underutilized accelerator interconnect bandwidth to assist host-to-accelerator data movement through adaptive multi-path transfer. Our experiments show that \system trains up to $1.68\times$ faster and MemFerry can train $1.52\times$ larger model compared to ZeRO-Offload on a single GPU, and increase training speed by at least $28.1%$ when scaling to data parallelism on 8 GPUs. We further extend the design to a Huawei CloudMatrix384 scale-Up node with up to 8 NPUs, and our ScaleUp-MemFerry reduces the end-to-end iteration time by up to $20.7%$ over DeepSpeed.
Mixture-of-Experts (MoE) layers replace the feed-forward block of a Transformer with E expert networks, and each token is routed to k of these experts. Under expert parallelism (EP) the experts are distributed across GPUs, and every MoE layer runs all-to-all collectives in the forward and backward passes to dispatch tokens to their experts and then combine the results. On a cluster with 8 AMD Instinct MI300X GPUs per node, these collectives can take 45% of the training step at EP32 with top-2 routing and 60% with top-6 routing. We find that early in pretraining routers have already learned to assign tokens to experts in correlated patterns, both within a layer and across layers. At top-2, 0.8% of the expert pairs in a layer are selected together by 42% of tokens, and the experts a token selects at one layer predict the experts it selects at the next layer. We use these correlations to keep more token--expert assignments on the token's own GPU, which reduces communication across GPUs and across nodes. Correlated expert placement puts experts that are often selected together on the same GPU. Combined with a dispatcher that sends each token to each GPU once, it removes up to 58% of dispatched rows. Token shuffling applies when sequence parallelism shards tokens across the EP group. It moves each token to the GPU predicted to hold its next-layer experts during the reduce-scatter that follows attention. On one node this raises the share of token--expert assignments served on the token's GPU from 12.5% to 59%. In Megatron-LM, across EP degrees from 8 to 64 with top-2 and top-6 routing, the two methods reduce all-to-all time by 1.16-2.63X and end-to-end step time by up to 1.41X. Neither method changes the models' underlying routing decisions or expert parameters.
展开完整摘要收起摘要↓
Mixture-of-Experts (MoE) layers replace the feed-forward block of a Transformer with E expert networks, and each token is routed to k of these experts. Under expert parallelism (EP) the experts are distributed across GPUs, and every MoE layer runs all-to-all collectives in the forward and backward passes to dispatch tokens to their experts and then combine the results. On a cluster with 8 AMD Instinct MI300X GPUs per node, these collectives can take 45% of the training step at EP32 with top-2 routing and 60% with top-6 routing. We find that early in pretraining routers have already learned to assign tokens to experts in correlated patterns, both within a layer and across layers. At top-2, 0.8% of the expert pairs in a layer are selected together by 42% of tokens, and the experts a token selects at one layer predict the experts it selects at the next layer. We use these correlations to keep more token--expert assignments on the token's own GPU, which reduces communication across GPUs and across nodes. Correlated expert placement puts experts that are often selected together on the same GPU. Combined with a dispatcher that sends each token to each GPU once, it removes up to 58% of dispatched rows. Token shuffling applies when sequence parallelism shards tokens across the EP group. It moves each token to the GPU predicted to hold its next-layer experts during the reduce-scatter that follows attention. On one node this raises the share of token--expert assignments served on the token's GPU from 12.5% to 59%. In Megatron-LM, across EP degrees from 8 to 64 with top-2 and top-6 routing, the two methods reduce all-to-all time by 1.16-2.63X and end-to-end step time by up to 1.41X. Neither method changes the models' underlying routing decisions or expert parameters.
Distributed deep learning relies on data, pipeline, tensor, and hybrid parallelism, yet fault-tolerance mechanisms are typically evaluated only on the architecture for which they were designed. This leaves practitioners with little guidance when choosing mechanisms across architectures. FailBench provides a unified evaluation harness covering seven distributed training architectures, eight crash-fault-tolerance mechanisms and a no-FT baseline, and single, concurrent, and cascading fail-stop failures. We evaluate 142 (architecture, mechanism, trace) combinations on an 8xV100 cluster, with per-rank checkpoint states ranging from 205 MB to 2.7 GB. Three findings emerge. First, no mechanism is universally best: on A2, disk checkpointing has the lowest steady-state overhead (0.5%), in-memory replication restores fastest (~17 ms), and just-in-time checkpointing avoids periodic steady-state checkpoint cost but incurs ~0.9 s upon failure. When process-group re-formation takes seconds, mechanisms differ more in runtime overhead than restore speed. Second, gossip training increases sample throughput by 17.9% after losing a worker, yet shows no detectable improvement in loss progress over a matched no-fault baseline. Third, mechanism cost depends strongly on architecture: in-memory replication overhead ranges from 3.7% to 176%. We translate these findings into a decision framework for selecting fault-tolerance mechanisms and release FailBench as an open artifact.
展开完整摘要收起摘要↓
Distributed deep learning relies on data, pipeline, tensor, and hybrid parallelism, yet fault-tolerance mechanisms are typically evaluated only on the architecture for which they were designed. This leaves practitioners with little guidance when choosing mechanisms across architectures. FailBench provides a unified evaluation harness covering seven distributed training architectures, eight crash-fault-tolerance mechanisms and a no-FT baseline, and single, concurrent, and cascading fail-stop failures. We evaluate 142 (architecture, mechanism, trace) combinations on an 8xV100 cluster, with per-rank checkpoint states ranging from 205 MB to 2.7 GB. Three findings emerge. First, no mechanism is universally best: on A2, disk checkpointing has the lowest steady-state overhead (0.5%), in-memory replication restores fastest (~17 ms), and just-in-time checkpointing avoids periodic steady-state checkpoint cost but incurs ~0.9 s upon failure. When process-group re-formation takes seconds, mechanisms differ more in runtime overhead than restore speed. Second, gossip training increases sample throughput by 17.9% after losing a worker, yet shows no detectable improvement in loss progress over a matched no-fault baseline. Third, mechanism cost depends strongly on architecture: in-memory replication overhead ranges from 3.7% to 176%. We translate these findings into a decision framework for selecting fault-tolerance mechanisms and release FailBench as an open artifact.
作者Hyungyo Kim, Nicholas Satchanov, Hrishi Shah, Gaohan Ye, Jiaqi Lou, Robert Walkup, Shweta Salaria, I-Hsin Chung, Hubertus Franke, Seetharami Seelam, Apoorve Mohan, Nam Sung Kim
TRANSIT is a transparent scale-in framework to enable multi-node model training on fewer GPUs while maintaining training efficiency by transparently leveraging CPU DRAM as an extension of GPU memory during distributed training. It achieves this through a user-space interposition layer, requiring no modifications to the application, training framework, cluster scheduler, device driver, or operating system. Furthermore, TRANSIT achieves higher efficiency by leveraging a zero-copy data path for CPU-GPU transfers. We evaluate TRANSIT on dense and MoE models across scales up to 64 NVIDIA H100 GPUs and multiple parallelism configurations over a RoCE network. Our evaluation shows that TRANSIT can: (a) outperform state-of-the-art framework-managed offloading techniques, achieving up to 68%, 59%, and 42% higher per-GPU throughput than TorchTitan, ZeRO-Offload, and ZeRO-Infinity, respectively, (b) enables training with 50% fewer GPUs while maintaining over 90% of baseline per-GPU throughput, (c) lower per-node network traffic by up to 33%, and (d) improve per-GPU throughput by up to 35% in communication-bound settings.
展开完整摘要收起摘要↓
TRANSIT is a transparent scale-in framework to enable multi-node model training on fewer GPUs while maintaining training efficiency by transparently leveraging CPU DRAM as an extension of GPU memory during distributed training. It achieves this through a user-space interposition layer, requiring no modifications to the application, training framework, cluster scheduler, device driver, or operating system. Furthermore, TRANSIT achieves higher efficiency by leveraging a zero-copy data path for CPU-GPU transfers. We evaluate TRANSIT on dense and MoE models across scales up to 64 NVIDIA H100 GPUs and multiple parallelism configurations over a RoCE network. Our evaluation shows that TRANSIT can: (a) outperform state-of-the-art framework-managed offloading techniques, achieving up to 68%, 59%, and 42% higher per-GPU throughput than TorchTitan, ZeRO-Offload, and ZeRO-Infinity, respectively, (b) enables training with 50% fewer GPUs while maintaining over 90% of baseline per-GPU throughput, (c) lower per-node network traffic by up to 33%, and (d) improve per-GPU throughput by up to 35% in communication-bound settings.
Reinforcement learning (RL) for post-training large language models (LLMs) incurs substantial computation and memory overhead during rollout generation, which motivates low-precision rollout for efficient RL training. However, existing FP4 RL methods suffer from a key limitation: they primarily optimize quantization accuracy on the training and rollout paths independently rather than directly reducing the discrepancy between the two quantized execution paths. In this work, we propose TRACE (Train-Rollout Quantization Alignment via Compact GuidancE), an FP4 quantization framework for RL training of Mixture-of-Experts (MoE) language models that addresses the limitation of existing FP4 RL methods. TRACE incorporates rollout-guided quantization-aware training that uses rollout-side quantization outcomes to guide training-side FP4 rounding decisions, directly reducing train-rollout discrepancy. Moreover, TRACE adopts an efficient quantization-information caching scheme that selectively retains mantissa and scale information from deeper layers to reduce the storage and communication overhead introduced by rollout guidance. We evaluate TRACE on four large-scale MoE language models across reasoning, coding, and long-horizon RL tasks. Our results demonstrate that TRACE enables joint FP4 weight/activation and FP4 KV-cache rollout with RL performance comparable to BF16 rollout, while achieving up to 5.4xrollout speedup and strong final FP4 performance compared with post-hoc FP4 quantization of BF16-trained policies.
展开完整摘要收起摘要↓
Reinforcement learning (RL) for post-training large language models (LLMs) incurs substantial computation and memory overhead during rollout generation, which motivates low-precision rollout for efficient RL training. However, existing FP4 RL methods suffer from a key limitation: they primarily optimize quantization accuracy on the training and rollout paths independently rather than directly reducing the discrepancy between the two quantized execution paths. In this work, we propose TRACE (Train-Rollout Quantization Alignment via Compact GuidancE), an FP4 quantization framework for RL training of Mixture-of-Experts (MoE) language models that addresses the limitation of existing FP4 RL methods. TRACE incorporates rollout-guided quantization-aware training that uses rollout-side quantization outcomes to guide training-side FP4 rounding decisions, directly reducing train-rollout discrepancy. Moreover, TRACE adopts an efficient quantization-information caching scheme that selectively retains mantissa and scale information from deeper layers to reduce the storage and communication overhead introduced by rollout guidance. We evaluate TRACE on four large-scale MoE language models across reasoning, coding, and long-horizon RL tasks. Our results demonstrate that TRACE enables joint FP4 weight/activation and FP4 KV-cache rollout with RL performance comparable to BF16 rollout, while achieving up to 5.4xrollout speedup and strong final FP4 performance compared with post-hoc FP4 quantization of BF16-trained policies.
Distributed training and rollout generation often use different tensor layouts, requiring model weights to be resharded across distinct process groups. This M-to-N redistribution is not directly expressed by standard collectives. Flat direct sends duplicate traffic across destination replicas, while gather-then-broadcast concentrates network injection at one root and transfers data that destinations do not need. For DeepSeek-V3 on 256 GPUs, legacy all-gather plus broadcast accounts for 29.4% of the reported reinforcement-learning step time. We present NCCL M2N, a layout- and topology-aware collective primitive for distributed tensor resharding. Given source and destination meshes and placements, it derives the required transfer regions and a global communication schedule. Its hierarchical route balances source contributions across eligible destination ranks, forwards one copy between destination NVLink domains, and completes local replication over NVLink. Network transfer and local replication overlap, avoiding replica-multiplied source egress. An aggregate data-movement model captures the limits of both stages. We evaluate NCCL M2N on up to 256 GB200 GPUs in an NVL72 cluster with NDR InfiniBand. A single FFN-MoE layer transfer achieves up to 7.9x speedup over flat direct sends (9.8 ms versus 77.3 ms). In a separate 256-GPU DeepSeek-V3 NeMo-RL experiment, NCCL M2N reduces reported weight-sync time from 5.78 s to 2.77 s, a 2.09x speedup over legacy all-gather plus broadcast, and reduces step time by 12.7%.
展开完整摘要收起摘要↓
Distributed training and rollout generation often use different tensor layouts, requiring model weights to be resharded across distinct process groups. This M-to-N redistribution is not directly expressed by standard collectives. Flat direct sends duplicate traffic across destination replicas, while gather-then-broadcast concentrates network injection at one root and transfers data that destinations do not need. For DeepSeek-V3 on 256 GPUs, legacy all-gather plus broadcast accounts for 29.4% of the reported reinforcement-learning step time. We present NCCL M2N, a layout- and topology-aware collective primitive for distributed tensor resharding. Given source and destination meshes and placements, it derives the required transfer regions and a global communication schedule. Its hierarchical route balances source contributions across eligible destination ranks, forwards one copy between destination NVLink domains, and completes local replication over NVLink. Network transfer and local replication overlap, avoiding replica-multiplied source egress. An aggregate data-movement model captures the limits of both stages. We evaluate NCCL M2N on up to 256 GB200 GPUs in an NVL72 cluster with NDR InfiniBand. A single FFN-MoE layer transfer achieves up to 7.9x speedup over flat direct sends (9.8 ms versus 77.3 ms). In a separate 256-GPU DeepSeek-V3 NeMo-RL experiment, NCCL M2N reduces reported weight-sync time from 5.78 s to 2.77 s, a 2.09x speedup over legacy all-gather plus broadcast, and reduces step time by 12.7%.
作者Arnab Kanti Tarafder, Jaume Guasch-Martí, Gokcen Kestor, Jie Ren
As Mixture-of-Experts (MoE) models scale toward hundreds of experts and higher top-$k$ routing, memory efficiency in distributed training becomes a critical bottleneck. Peak memory is dominated by the MoE block, not attention: every intermediate buffer in the MoE dispatch pipeline is individually scaled by top-k routing. The standard all-to-all dispatcher sends all routed tokens in a single collective step, requiring the full top-$k$-expanded buffer to be constructed at once. In this work, we propose RelayMoE, a ring-based MoE execution model that computes locally as expert weights or tokens circulate, avoiding full top-$k$-expanded dispatch buffers. RelayMoE selects between expert and token routing according to communication volume and overlaps transfers with computation. The ring structure naturally supports memory-efficient MoE recomputation during backward: each hop reconstructs expert intermediates, uses them to compute gradients, and releases them before the next hop. The saved memory supports longer sequences and larger batches, or retains more attention activations to reduce attention recomputation and improve training throughput. We evaluate RelayMoE on 30B$-$57B production MoE models and varied expert configurations. In single-layer MoE experiments, RelayMoE achieves a $2\times$ average speedup over Megatron-LM. In full-model training under the same GPU memory budget, it improves throughput by up to $2.02\times$ and extends the largest tested trainable sequence length by up to $2.85\times$.
展开完整摘要收起摘要↓
As Mixture-of-Experts (MoE) models scale toward hundreds of experts and higher top-$k$ routing, memory efficiency in distributed training becomes a critical bottleneck. Peak memory is dominated by the MoE block, not attention: every intermediate buffer in the MoE dispatch pipeline is individually scaled by top-k routing. The standard all-to-all dispatcher sends all routed tokens in a single collective step, requiring the full top-$k$-expanded buffer to be constructed at once. In this work, we propose RelayMoE, a ring-based MoE execution model that computes locally as expert weights or tokens circulate, avoiding full top-$k$-expanded dispatch buffers. RelayMoE selects between expert and token routing according to communication volume and overlaps transfers with computation. The ring structure naturally supports memory-efficient MoE recomputation during backward: each hop reconstructs expert intermediates, uses them to compute gradients, and releases them before the next hop. The saved memory supports longer sequences and larger batches, or retains more attention activations to reduce attention recomputation and improve training throughput. We evaluate RelayMoE on 30B$-$57B production MoE models and varied expert configurations. In single-layer MoE experiments, RelayMoE achieves a $2\times$ average speedup over Megatron-LM. In full-model training under the same GPU memory budget, it improves throughput by up to $2.02\times$ and extends the largest tested trainable sequence length by up to $2.85\times$.
作者Guilherme Silva, Pedro Silva, Gladston Moreira, Eduardo Luz
Wearable and bedside electrocardiogram (ECG) monitors must adapt to patient-specific morphology to maintain arrhythmia detection accuracy across users, yet personalization is typically performed offline and cannot account for individual physiology, electrode placement, or recording drift. On-device adaptation by backpropagation is expensive for microcontroller-class medical devices because it requires an optimizer state, repeated backward passes through convolutional layers, and labeled arrhythmic beats that may not be available at deployment time. This letter proposes prototype-only head adaptation as a compact personalization primitive for TinyML ECG systems. A one-dimensional convolutional neural network (1-D CNN; 1,314 parameters and 72.6k multiply-accumulate operations per beat) is trained offline on the MIT-BIH Arrhythmia Database under an inter-patient protocol, frozen as a feature extractor, and exported to a PSoC 6 microcontroller. Patient-specific adaptation then reduces to computing closed-form class means in a 32-dimensional embedding space, requiring no convolutional backward pass, no iterative optimization, and only one forward pass per support beat. Prototype adaptation improves inter-patient macro-F1 from 0.635/0.639/0.646 to 0.731/0.771/0.797 at 1/5/10-shot, outperforming linear stochastic-gradient-descent (SGD) head fine-tuning at every shot count for the target tiny backbone. On-device replay over 18 one-shot episodes on a PSoC 6 Cortex-M4F matches the host macro-F1 for the prototype head (0.798), with 11.39 ms per beat, 5.2 KB flash, and 22.2 KB SRAM. A restricted variant that updates only the normal-class prototype from passively buffered sinus beats yields a consistent +0.05 macro-F1 gain, reducing the annotation burden during initial
展开完整摘要收起摘要↓
Wearable and bedside electrocardiogram (ECG) monitors must adapt to patient-specific morphology to maintain arrhythmia detection accuracy across users, yet personalization is typically performed offline and cannot account for individual physiology, electrode placement, or recording drift. On-device adaptation by backpropagation is expensive for microcontroller-class medical devices because it requires an optimizer state, repeated backward passes through convolutional layers, and labeled arrhythmic beats that may not be available at deployment time. This letter proposes prototype-only head adaptation as a compact personalization primitive for TinyML ECG systems. A one-dimensional convolutional neural network (1-D CNN; 1,314 parameters and 72.6k multiply-accumulate operations per beat) is trained offline on the MIT-BIH Arrhythmia Database under an inter-patient protocol, frozen as a feature extractor, and exported to a PSoC 6 microcontroller. Patient-specific adaptation then reduces to computing closed-form class means in a 32-dimensional embedding space, requiring no convolutional backward pass, no iterative optimization, and only one forward pass per support beat. Prototype adaptation improves inter-patient macro-F1 from 0.635/0.639/0.646 to 0.731/0.771/0.797 at 1/5/10-shot, outperforming linear stochastic-gradient-descent (SGD) head fine-tuning at every shot count for the target tiny backbone. On-device replay over 18 one-shot episodes on a PSoC 6 Cortex-M4F matches the host macro-F1 for the prototype head (0.798), with 11.39 ms per beat, 5.2 KB flash, and 22.2 KB SRAM. A restricted variant that updates only the normal-class prototype from passively buffered sinus beats yields a consistent +0.05 macro-F1 gain, reducing the annotation burden during initial
Parameter-efficient fine-tuning (PEFT) reduces the cost of adapting foundation models by focusing training on a small parameter subset. Complementary to this idea, we introduce RoSA (Rotational Sparse Adaptation), which narrows adaptation to a subset of layers at a time. RoSA freezes lower layers close to the input throughout training and rotates a trainable block over later layers, progressively increasing the number of frozen layers close to the input. This design reduces optimizer-state memory, shortens backpropagation, and even forward propagation if activations at the last frozen layer are cached. Because RoSA is orthogonal to the choice of trainable parameterization, it can be combined with PEFT methods or sparse optimizers within each active block. Experiments across multiple LLM architectures and tasks show that RoSA reduces peak memory while maintaining strong fine-tuning performance.
展开完整摘要收起摘要↓
Parameter-efficient fine-tuning (PEFT) reduces the cost of adapting foundation models by focusing training on a small parameter subset. Complementary to this idea, we introduce RoSA (Rotational Sparse Adaptation), which narrows adaptation to a subset of layers at a time. RoSA freezes lower layers close to the input throughout training and rotates a trainable block over later layers, progressively increasing the number of frozen layers close to the input. This design reduces optimizer-state memory, shortens backpropagation, and even forward propagation if activations at the last frozen layer are cached. Because RoSA is orthogonal to the choice of trainable parameterization, it can be combined with PEFT methods or sparse optimizers within each active block. Experiments across multiple LLM architectures and tasks show that RoSA reduces peak memory while maintaining strong fine-tuning performance.
Mixture-of-Experts (MoE) has been widely adopted in recent large language model (LLM) architectures. However, scaling up MoE in LLM training introduces system-level challenges on training, where non-uniform token routing can lead to highly imbalanced workloads across experts and devices, further destabilizing the training process. With trillion-scale LLMs, imbalanced expert workloads further amplify the resource cost of MoE training, resulting in degraded training efficiency and hardware utilization for underloaded experts, while hot experts require additional resources to accommodate excessive workloads. Recent studies address imbalanced MoE training through intricate parallelism strategies or resource reallocation. However, these system-level approaches often introduce additional resource requirements and considerable orchestration complexity, which become increasingly difficult to afford when training trillion-parameter LLMs under constrained computational resources. This work introduces CIPHER-MoE, which mitigates MoE workload imbalance while keeping the router's token-side Top-K selection unchanged. CIPHER-MoE applies affinity-aware Expert-to-Token filtering with explicit capacity control to reduce hotspot expert workloads without additional hardware resources or complex runtime design. The proposed method has been evaluated on large-scale MoE models, including DeepSeek-V4-Pro, showing up to 64.9 percentage points Top-1 expert workload reduction and 1.10$\times$-1.94$\times$ training acceleration, while preserving the training quality. The source code will be released soon.
展开完整摘要收起摘要↓
Mixture-of-Experts (MoE) has been widely adopted in recent large language model (LLM) architectures. However, scaling up MoE in LLM training introduces system-level challenges on training, where non-uniform token routing can lead to highly imbalanced workloads across experts and devices, further destabilizing the training process. With trillion-scale LLMs, imbalanced expert workloads further amplify the resource cost of MoE training, resulting in degraded training efficiency and hardware utilization for underloaded experts, while hot experts require additional resources to accommodate excessive workloads. Recent studies address imbalanced MoE training through intricate parallelism strategies or resource reallocation. However, these system-level approaches often introduce additional resource requirements and considerable orchestration complexity, which become increasingly difficult to afford when training trillion-parameter LLMs under constrained computational resources. This work introduces CIPHER-MoE, which mitigates MoE workload imbalance while keeping the router's token-side Top-K selection unchanged. CIPHER-MoE applies affinity-aware Expert-to-Token filtering with explicit capacity control to reduce hotspot expert workloads without additional hardware resources or complex runtime design. The proposed method has been evaluated on large-scale MoE models, including DeepSeek-V4-Pro, showing up to 64.9 percentage points Top-1 expert workload reduction and 1.10$\times$-1.94$\times$ training acceleration, while preserving the training quality. The source code will be released soon.
作者Zishan Shao, Liang Tian, Georgiy Zemlevskiy, Kangning Cui, Lixun Zhang, Yixiao Wang, Ting Jiang, Jinhee Kim, Yixuan Chen, Rui-Feng Wang, Fan Yang, Hai Li, Yiran Chen
Different low-rank compression methods can produce compressed LLMs that respond differently to the same post-compression recovery procedure, and relative advantages observed between methods at the compression endpoint may shrink, grow, or even reverse after recovery. We ask whether this recovery heterogeneity reflects functional structure beyond scalar loss evolution, and how that structure evolves throughout recovery. Our results establish that this heterogeneity reflects a reproducible compression-induced functional structure, which we formalize as recovery pressure. To characterize this structure consistently throughout recovery, we develop a standardized functional characterization within each backbone that is applicable across heterogeneous low-rank methods. The primary backward characterization reveals reproducible module-wise structure across independent probes, while a complementary forward-only characterization recovers related structure without loss or backpropagation. We further find that recovery pressure measured at the endpoint is associated with subsequent recovery response; during recovery, its module-wise structure is reorganized non-uniformly, and localized updates induce distributed responses beyond directly updated modules. Further evidence indicates that tracking this evolving structure provides a complementary functional view of recovery progress alongside scalar loss.
展开完整摘要收起摘要↓
Different low-rank compression methods can produce compressed LLMs that respond differently to the same post-compression recovery procedure, and relative advantages observed between methods at the compression endpoint may shrink, grow, or even reverse after recovery. We ask whether this recovery heterogeneity reflects functional structure beyond scalar loss evolution, and how that structure evolves throughout recovery. Our results establish that this heterogeneity reflects a reproducible compression-induced functional structure, which we formalize as recovery pressure. To characterize this structure consistently throughout recovery, we develop a standardized functional characterization within each backbone that is applicable across heterogeneous low-rank methods. The primary backward characterization reveals reproducible module-wise structure across independent probes, while a complementary forward-only characterization recovers related structure without loss or backpropagation. We further find that recovery pressure measured at the endpoint is associated with subsequent recovery response; during recovery, its module-wise structure is reorganized non-uniformly, and localized updates induce distributed responses beyond directly updated modules. Further evidence indicates that tracking this evolving structure provides a complementary functional view of recovery progress alongside scalar loss.
Training large language models (LLMs) entails a fundamental trade-off: memory-efficient optimizers such as Adam discard cross-parameter curvature, whereas full-curvature methods such as SOAP can accelerate convergence at prohibitive memory costs. We introduce Clean, a memory-efficient and full-curvature optimizer designed to resolve this bottleneck. Clean leverages the randomized Nystrom method to accurately approximate the left and right preconditioners in SOAP, and to reduce the optimizer's memory complexity from quadratic to linear in terms of model dimensions. We subsequently reintegrate the off-subspace components to capture curvature information beyond the low-rank approximation, preserving rich curvature at minimal memory cost. We further propose Q-Clean, a low-precision variant that aggressively compresses optimizer states. Q-Clean reduces optimizer memory consumption by over 50% compared to Muon when pre-training a LLaMA-1.3B architecture, all while maintaining strong and competitive predictive performance. Notably, Clean operates with a smaller optimizer-state footprint than standard AdamW while reaching AdamW's final performance 26% faster in wall-clock time. Furthermore, our methods uniquely enable the pre-training of a 13B-parameter model on a single 80GB GPU, providing a scalable, efficient, and accessible approach to large-scale model optimization.
展开完整摘要收起摘要↓
Training large language models (LLMs) entails a fundamental trade-off: memory-efficient optimizers such as Adam discard cross-parameter curvature, whereas full-curvature methods such as SOAP can accelerate convergence at prohibitive memory costs. We introduce Clean, a memory-efficient and full-curvature optimizer designed to resolve this bottleneck. Clean leverages the randomized Nystrom method to accurately approximate the left and right preconditioners in SOAP, and to reduce the optimizer's memory complexity from quadratic to linear in terms of model dimensions. We subsequently reintegrate the off-subspace components to capture curvature information beyond the low-rank approximation, preserving rich curvature at minimal memory cost. We further propose Q-Clean, a low-precision variant that aggressively compresses optimizer states. Q-Clean reduces optimizer memory consumption by over 50% compared to Muon when pre-training a LLaMA-1.3B architecture, all while maintaining strong and competitive predictive performance. Notably, Clean operates with a smaller optimizer-state footprint than standard AdamW while reaching AdamW's final performance 26% faster in wall-clock time. Furthermore, our methods uniquely enable the pre-training of a 13B-parameter model on a single 80GB GPU, providing a scalable, efficient, and accessible approach to large-scale model optimization.
The Muon optimizer derives its update rule for hidden linear layers by solving a local linearization of the loss penalized by the spectral norm, motivated by an RMS-stability argument for dense linear layers. Standard Muon implementations, however, exclude the input (embedding table) and output (language model head) layers from this principled treatment, for which they use AdamW instead. We present MuonIO, a single Muon-style update for both of these layers. For the language model head $\mathbf{L} \in \mathbb{R}^{V \times d}$, we motivate the use of the $2\to\infty$ operator norm, due to the Lipschitz continuity of the softmax output geometry, while for the embedding table $\mathbf{E} \in \mathbb{R}^{d \times V}$, we draw on the $1 \to 2$ operator norm, based on the one-hot input geometry identified by Bernstein & Newhouse (2025). The identity $\lVert\mathbf{L}\rVert_{2\to\infty}=\lVert\mathbf{L}^\top\rVert_{1\to2}$ then puts both matrices in the same vocabulary-oriented geometry: MuonIO applies a single normalized-vector rule, which appears as column normalization for $\mathbf{E}$ and row normalization for $\mathbf{L}$. Empirical evaluations demonstrate the effectiveness of our approach, with MuonIO reducing I/O optimizer state memory by 50% and I/O update FLOPs by $\sim$46% compared to Muon for 1B LLaMA pretraining on C4, while also improving validation perplexity.
展开完整摘要收起摘要↓
The Muon optimizer derives its update rule for hidden linear layers by solving a local linearization of the loss penalized by the spectral norm, motivated by an RMS-stability argument for dense linear layers. Standard Muon implementations, however, exclude the input (embedding table) and output (language model head) layers from this principled treatment, for which they use AdamW instead. We present MuonIO, a single Muon-style update for both of these layers. For the language model head $\mathbf{L} \in \mathbb{R}^{V \times d}$, we motivate the use of the $2\to\infty$ operator norm, due to the Lipschitz continuity of the softmax output geometry, while for the embedding table $\mathbf{E} \in \mathbb{R}^{d \times V}$, we draw on the $1 \to 2$ operator norm, based on the one-hot input geometry identified by Bernstein & Newhouse (2025). The identity $\lVert\mathbf{L}\rVert_{2\to\infty}=\lVert\mathbf{L}^\top\rVert_{1\to2}$ then puts both matrices in the same vocabulary-oriented geometry: MuonIO applies a single normalized-vector rule, which appears as column normalization for $\mathbf{E}$ and row normalization for $\mathbf{L}$. Empirical evaluations demonstrate the effectiveness of our approach, with MuonIO reducing I/O optimizer state memory by 50% and I/O update FLOPs by $\sim$46% compared to Muon for 1B LLaMA pretraining on C4, while also improving validation perplexity.
Irregular All-to-All communication is a major bottleneck in expert-parallel Mixture-of-Experts (MoE) models. Even with fixed expert routing and placement, uneven utilization of parallel network Rails and incast can limit communication performance. We present RailWave, a phase-adaptive communication layer built on DeepEP that addresses these bottlenecks below the routing layer through spatial and temporal traffic shaping. RailBalance redistributes source traffic across eligible Rails using source-local information, while a reusable, topology-derived permutation schedule limits concurrent senders per receiver without rebuilding demand-dependent schedules for each communication phase. A lightweight calibrated selector chooses an execution path according to each phase's traffic characteristics and offline profiling results. On training-derived communication workloads from the 106B GLM-4.5-Air model, RailWave delivers up to 5.84x speedup on H800 and 4.36x on H20 over Native. Code is available at https://github.com/CyberSecurityErial/RailWave-EP.
展开完整摘要收起摘要↓
Irregular All-to-All communication is a major bottleneck in expert-parallel Mixture-of-Experts (MoE) models. Even with fixed expert routing and placement, uneven utilization of parallel network Rails and incast can limit communication performance. We present RailWave, a phase-adaptive communication layer built on DeepEP that addresses these bottlenecks below the routing layer through spatial and temporal traffic shaping. RailBalance redistributes source traffic across eligible Rails using source-local information, while a reusable, topology-derived permutation schedule limits concurrent senders per receiver without rebuilding demand-dependent schedules for each communication phase. A lightweight calibrated selector chooses an execution path according to each phase's traffic characteristics and offline profiling results. On training-derived communication workloads from the 106B GLM-4.5-Air model, RailWave delivers up to 5.84x speedup on H800 and 4.36x on H20 over Native. Code is available at https://github.com/CyberSecurityErial/RailWave-EP.
作者Arash Lagzian, Paniz Halvachi, Junming Zhang, Zhouhan Lin, Dianbo Liu
Muon improves large-scale training by applying a spectral-norm steepest-descent update to matrix parameters, but practical models also contain parameter blocks that do not fit dense-matrix geometry. One important case is the tied vocabulary table, which appears in language models and other token generators and can receive multiple structurally different gradient sources, from sparse input lookups to dense output-classifier updates. In the reference recipe these blocks are handed to an auxiliary AdamW optimizer, which restores second-moment state and updates the aliased table as a generic tensor. We propose AF-Muon, an AdamW-free extension of Muon that keeps the Muon matrix update for hidden weight matrices while using a support-aware finite-cap linear minimization oracle for tied vocabulary tables and an RMS-normalized update for one-dimensional auxiliary parameters. AF-Muon therefore trains every parameter class with a single first-moment buffer and no second-moment state, saving around 20% optimizer-state memory relative to Hybrid Muon in our benchmark. Across nine tied-token settings - decoder-only language models from 124M to 1B parameters, a fully shared T5-style encoder-decoder, and ImageGPT-style image-token, protein, and sparse-MoE variants, spanning text, image, and protein-sequence data - AF-Muon improves mean validation loss and perplexity over both Hybrid Muon and a SCION-style Sign endpoint. Long-horizon runs and hyperparameter sensitivity studies confirm the gain is robust, and identical-momentum diagnostics attribute it to the finite cap, which preserves more within-row magnitude than Sign while bounding the coordinate concentration of row-RMS. These results identify tied vocabulary tables as a distinct optimizer geometry and yield a robust AdamW-free Muon variant across models, modalities, and architectures, with about 1% step-time overhead in matched training.
展开完整摘要收起摘要↓
Muon improves large-scale training by applying a spectral-norm steepest-descent update to matrix parameters, but practical models also contain parameter blocks that do not fit dense-matrix geometry. One important case is the tied vocabulary table, which appears in language models and other token generators and can receive multiple structurally different gradient sources, from sparse input lookups to dense output-classifier updates. In the reference recipe these blocks are handed to an auxiliary AdamW optimizer, which restores second-moment state and updates the aliased table as a generic tensor. We propose AF-Muon, an AdamW-free extension of Muon that keeps the Muon matrix update for hidden weight matrices while using a support-aware finite-cap linear minimization oracle for tied vocabulary tables and an RMS-normalized update for one-dimensional auxiliary parameters. AF-Muon therefore trains every parameter class with a single first-moment buffer and no second-moment state, saving around 20% optimizer-state memory relative to Hybrid Muon in our benchmark. Across nine tied-token settings - decoder-only language models from 124M to 1B parameters, a fully shared T5-style encoder-decoder, and ImageGPT-style image-token, protein, and sparse-MoE variants, spanning text, image, and protein-sequence data - AF-Muon improves mean validation loss and perplexity over both Hybrid Muon and a SCION-style Sign endpoint. Long-horizon runs and hyperparameter sensitivity studies confirm the gain is robust, and identical-momentum diagnostics attribute it to the finite cap, which preserves more within-row magnitude than Sign while bounding the coordinate concentration of row-RMS. These results identify tied vocabulary tables as a distinct optimizer geometry and yield a robust AdamW-free Muon variant across models, modalities, and architectures, with about 1% step-time overhead in matched training.
Federated training of foundation models is constrained by client memory and communication costs. LoRA-based methods reduce these costs through low-rank adapters, but their fixed rank budget can limit adaptation. Gradient low-rank optimization offers greater flexibility, yet independently chosen client subspaces create a problem we term subspace fragmentation: local projections interact with data heterogeneity to bias aggregated directions, while aggregation can increase update rank and communication cost. Thus, accurate local gradient compression need not preserve global descent. We propose FedLore, which shares a low-rank optimization basis within each round and refreshes it across rounds. The shared basis enables exact aggregation in low-rank coordinates and eliminates the identified projection bias. Subspace refresh allows the accumulated model update to exceed the per-round rank budget. We characterize the aggregation bias and establish an $O(T^{-1/2})$ stationarity bound for the projected-SGD variant under a global-gradient coverage condition and standard smoothness and variance assumptions, with bounded gradient heterogeneity. Experiments on vision and language tasks, including federated pre-training, show that FedLore outperforms the evaluated low-rank adapter baselines and matches or exceeds full-parameter training, while reducing communication and optimizer-state memory.
展开完整摘要收起摘要↓
Federated training of foundation models is constrained by client memory and communication costs. LoRA-based methods reduce these costs through low-rank adapters, but their fixed rank budget can limit adaptation. Gradient low-rank optimization offers greater flexibility, yet independently chosen client subspaces create a problem we term subspace fragmentation: local projections interact with data heterogeneity to bias aggregated directions, while aggregation can increase update rank and communication cost. Thus, accurate local gradient compression need not preserve global descent. We propose FedLore, which shares a low-rank optimization basis within each round and refreshes it across rounds. The shared basis enables exact aggregation in low-rank coordinates and eliminates the identified projection bias. Subspace refresh allows the accumulated model update to exceed the per-round rank budget. We characterize the aggregation bias and establish an $O(T^{-1/2})$ stationarity bound for the projected-SGD variant under a global-gradient coverage condition and standard smoothness and variance assumptions, with bounded gradient heterogeneity. Experiments on vision and language tasks, including federated pre-training, show that FedLore outperforms the evaluated low-rank adapter baselines and matches or exceeds full-parameter training, while reducing communication and optimizer-state memory.
作者Seungjun Lee, Ensieh Khazaei, Dimitrios Hatzinakos, Baturalp Buyukates, Sunwoo Lee
Federated learning (FL) is a communication-efficient distributed learning paradigm. However, client drift remains one of the most critical challenges, hindering the efficient training of a global model. In this study, we propose a novel latent information sharing scheme that directly mitigates data heterogeneity across clients. Our theoretical and empirical results show that sharing a small amount of hidden-layer activations significantly improves training efficiency while preserving convergence guarantees and data privacy. Furthermore, we compare our method with existing FL approaches designed to address client drift, including FedProx, SCAFFOLD, FedPVR, FedProto, and SplitFed, and demonstrate superior model accuracy under a fixed round budget without incurring excessive communication overhead. Overall, this work presents a promising new knowledge aggregation scheme and provides a comprehensive analysis of the impact of activation sharing on federated optimization.
展开完整摘要收起摘要↓
Federated learning (FL) is a communication-efficient distributed learning paradigm. However, client drift remains one of the most critical challenges, hindering the efficient training of a global model. In this study, we propose a novel latent information sharing scheme that directly mitigates data heterogeneity across clients. Our theoretical and empirical results show that sharing a small amount of hidden-layer activations significantly improves training efficiency while preserving convergence guarantees and data privacy. Furthermore, we compare our method with existing FL approaches designed to address client drift, including FedProx, SCAFFOLD, FedPVR, FedProto, and SplitFed, and demonstrate superior model accuracy under a fixed round budget without incurring excessive communication overhead. Overall, this work presents a promising new knowledge aggregation scheme and provides a comprehensive analysis of the impact of activation sharing on federated optimization.
Full-parameter fine-tuning of large language models (LLMs) incurs substantial optimizer state memory overhead, limiting the model sizes that fit on modern GPUs. Existing approaches either compress optimizer state, abandon first-order gradients, or change the update geometry while retaining dense state. The recently introduced Muon optimizer reduces optimizer memory through matrix-valued updates. Still, its geometry differs from AdamW and can lead to performance degradation when fine-tuning AdamW-pretrained models. To reduce optimizer memory without sacrificing accuracy or computational efficiency in LLM fine-tuning, we propose Ternary Absolute-max Column-wise One-sparse optimizer, or TACO, which follows Muon's operator-norm steepest-descent view but takes the geometric route further. TACO computes the exact steepest-descent direction under a dimension-normalized $1\to1$ operator norm by selecting the sign of the largest magnitude entry in each column of two-dimensional weight matrices. This retains first-order gradients while making optimizer state memory nearly negligible. Our practical TACO optimizer maintains only a small set of low precision gradient components per column, reducing persistent optimizer state by $174\times$ relative to AdamW8bit (from 27.7 GB to 0.16 GB) and peak training memory by $2.9\times$ (from 80.6 GB to 27.5 GB) on OPT-13B, while achieving comparable accuracy and runtime. TACO further enables full-parameter fine-tuning of 30-32B-parameter models on a single 80 GB H100 GPU across multiple model families and tasks.
展开完整摘要收起摘要↓
Full-parameter fine-tuning of large language models (LLMs) incurs substantial optimizer state memory overhead, limiting the model sizes that fit on modern GPUs. Existing approaches either compress optimizer state, abandon first-order gradients, or change the update geometry while retaining dense state. The recently introduced Muon optimizer reduces optimizer memory through matrix-valued updates. Still, its geometry differs from AdamW and can lead to performance degradation when fine-tuning AdamW-pretrained models. To reduce optimizer memory without sacrificing accuracy or computational efficiency in LLM fine-tuning, we propose Ternary Absolute-max Column-wise One-sparse optimizer, or TACO, which follows Muon's operator-norm steepest-descent view but takes the geometric route further. TACO computes the exact steepest-descent direction under a dimension-normalized $1\to1$ operator norm by selecting the sign of the largest magnitude entry in each column of two-dimensional weight matrices. This retains first-order gradients while making optimizer state memory nearly negligible. Our practical TACO optimizer maintains only a small set of low precision gradient components per column, reducing persistent optimizer state by $174\times$ relative to AdamW8bit (from 27.7 GB to 0.16 GB) and peak training memory by $2.9\times$ (from 80.6 GB to 27.5 GB) on OPT-13B, while achieving comparable accuracy and runtime. TACO further enables full-parameter fine-tuning of 30-32B-parameter models on a single 80 GB H100 GPU across multiple model families and tasks.
作者Peng Jin, Zihan Qiu, Zekun Wang, Bo Zheng, Yang Xu, Tian Xie, Xiao Li, Huaqing Zhang, Haoran Lian, Rui Men, Dayiheng Liu
Scaling Large Language Models (LLMs) via Mixture-of-Experts (MoE) enables massive parameter growth with nearly constant per-token computation. However, further scaling the parameter count requires increasingly sparse routing, where expert load imbalance becomes more severe. This imbalance reduces parameter utilization and training efficiency, and can undermine training stability, becoming a bottleneck to reliable scaling. In this work, we unify two representative auxiliary-loss-free methods as incomplete Proportional-Integral-Derivative (PID) controllers: DeepSeek's loss-free method acts as a fixed-step integral controller, while Kimi K3's Quantile Balancing functions as a generalized proportional controller. Building on this control perspective, we propose ID Balancing, an Integral-Derivative controller. It scales its integral term with load error and activates its derivative term only when imbalance worsens, enabling stronger corrections for large or worsening errors and smaller updates near balance. Evaluated across Top-$10$, Top-$5$, and Top-$3$ routing over $768$ experts, ID Balancing reduces worst-case backbone MaxVio and training-average backbone MinVio by over $50%$ and $12%$, respectively, relative to the best baselines in the Top-$3$ setting. When the total parameter count increases from $18.9$B to $69.9$B (Top-$10$-of-$768$), ID Balancing's worst-case backbone MaxVio remains nearly unchanged and is approximately $89.6%$ lower than that of the auxiliary-loss baseline. ID Balancing also maintains competitive language-modeling and downstream performance. The advantages of ID Balancing grow as sparsity increases, making it a promising solution for scaling larger, sparser MoE models.
展开完整摘要收起摘要↓
Scaling Large Language Models (LLMs) via Mixture-of-Experts (MoE) enables massive parameter growth with nearly constant per-token computation. However, further scaling the parameter count requires increasingly sparse routing, where expert load imbalance becomes more severe. This imbalance reduces parameter utilization and training efficiency, and can undermine training stability, becoming a bottleneck to reliable scaling. In this work, we unify two representative auxiliary-loss-free methods as incomplete Proportional-Integral-Derivative (PID) controllers: DeepSeek's loss-free method acts as a fixed-step integral controller, while Kimi K3's Quantile Balancing functions as a generalized proportional controller. Building on this control perspective, we propose ID Balancing, an Integral-Derivative controller. It scales its integral term with load error and activates its derivative term only when imbalance worsens, enabling stronger corrections for large or worsening errors and smaller updates near balance. Evaluated across Top-$10$, Top-$5$, and Top-$3$ routing over $768$ experts, ID Balancing reduces worst-case backbone MaxVio and training-average backbone MinVio by over $50%$ and $12%$, respectively, relative to the best baselines in the Top-$3$ setting. When the total parameter count increases from $18.9$B to $69.9$B (Top-$10$-of-$768$), ID Balancing's worst-case backbone MaxVio remains nearly unchanged and is approximately $89.6%$ lower than that of the auxiliary-loss baseline. ID Balancing also maintains competitive language-modeling and downstream performance. The advantages of ID Balancing grow as sparsity increases, making it a promising solution for scaling larger, sparser MoE models.
作者Fan Yang, Ying Zhou, Binglei Wang, Zhenjie Zhou, Jialong Li
AI clusters increasingly run large language model (LLM) inference and training on the same fabric. Prefill-decode (P-D) disaggregation creates key-value (KV) cache transfers between prefill and decode groups, whereas training collectives and all-to-all traffic benefit from near-uniform global connectivity. A static sparse topology can therefore be poorly matched to one of the two traffic patterns. We present Mode-Switching Expander (MoSE), a reconfigurable expander that treats topology design as a fixed-degree edge-allocation problem. MoSE reallocates the same sparse edge budget toward direct P-D connectivity in inference-heavy modes and restores a uniform random regular expander in training-heavy modes. We evaluate MoSE using a 1024-group flow-level topology model, shortest-path routing, and two mixed workloads. Across 20 seeds, MoSE reduces average and 95th-percentile (P95) load-aware KV communication cost by 90.8% and 91.9% relative to Static-Training in the inference-heavy mode. In the training-heavy mode, it reduces average and P95 training communication cost by 22.7% and 27.6% relative to stale Static-Inference. These results show that coarse-grained topology switching can support both workload modes without additional ports or routing changes.
展开完整摘要收起摘要↓
AI clusters increasingly run large language model (LLM) inference and training on the same fabric. Prefill-decode (P-D) disaggregation creates key-value (KV) cache transfers between prefill and decode groups, whereas training collectives and all-to-all traffic benefit from near-uniform global connectivity. A static sparse topology can therefore be poorly matched to one of the two traffic patterns. We present Mode-Switching Expander (MoSE), a reconfigurable expander that treats topology design as a fixed-degree edge-allocation problem. MoSE reallocates the same sparse edge budget toward direct P-D connectivity in inference-heavy modes and restores a uniform random regular expander in training-heavy modes. We evaluate MoSE using a 1024-group flow-level topology model, shortest-path routing, and two mixed workloads. Across 20 seeds, MoSE reduces average and 95th-percentile (P95) load-aware KV communication cost by 90.8% and 91.9% relative to Static-Training in the inference-heavy mode. In the training-heavy mode, it reduces average and P95 training communication cost by 22.7% and 27.6% relative to stale Static-Inference. These results show that coarse-grained topology switching can support both workload modes without additional ports or routing changes.
Hardware-operable failures (HOFs) interrupt large language model (LLM) training but permit recovery on the same hardware without reset, repair, or replacement. Existing recovery systems nevertheless reload checkpoints, recompute lost progress, and rebuild process state, idling GPUs that could otherwise continue training. We present Leto, a fault-tolerant training system that leverages surviving hardware to enable efficient in-place recovery. Our key insight is that the state needed to resume training can be retained or prepared outside the active training process while remaining on the same hardware. Leto retains the working model state and the reusable process state, and preinitializes the remaining state in a shadow trainer. We devise two-tier erasure protection and chunk-level transactional updates to keep the retained model state recoverable and consistent, and reclaim the shadow state when active training needs its GPU memory. Evaluation on 6- and 72-GPU NVIDIA A100 clusters shows that Leto recovers 3.6--6.5$\times$ faster than the best-performing checkpointing baselines and improves productive training time by up to 13.7 percentage points. Large-scale simulation shows over 95% productive training time on a 131,072-GPU cluster.
展开完整摘要收起摘要↓
Hardware-operable failures (HOFs) interrupt large language model (LLM) training but permit recovery on the same hardware without reset, repair, or replacement. Existing recovery systems nevertheless reload checkpoints, recompute lost progress, and rebuild process state, idling GPUs that could otherwise continue training. We present Leto, a fault-tolerant training system that leverages surviving hardware to enable efficient in-place recovery. Our key insight is that the state needed to resume training can be retained or prepared outside the active training process while remaining on the same hardware. Leto retains the working model state and the reusable process state, and preinitializes the remaining state in a shadow trainer. We devise two-tier erasure protection and chunk-level transactional updates to keep the retained model state recoverable and consistent, and reclaim the shadow state when active training needs its GPU memory. Evaluation on 6- and 72-GPU NVIDIA A100 clusters shows that Leto recovers 3.6--6.5$\times$ faster than the best-performing checkpointing baselines and improves productive training time by up to 13.7 percentage points. Large-scale simulation shows over 95% productive training time on a 131,072-GPU cluster.
Mixture-of-experts (MoE) megakernels fuse expert-parallel communication with expert computation. However, under fixed expert placement, routing skew creates GPU stragglers: overloaded GPUs determine layer latency while others sit idle. Replicating hot experts can shift work to underloaded GPUs, but dynamic replicas introduce additional work: replicas must receive expert weights to execute and, during training, their partial weight gradients must be reduced at the expert owners. We present MegaFlux, which makes expert replication a runtime decision and pipelines the communication induced by replication within persistent MoE execution. An on-device planner jointly selects replica locations and assigns tile-aligned token blocks under a per-GPU replica budget, leaving router outputs unchanged. The forward and backward megakernels realize pipelined expert replication: replicas begin computation as their required weights arrive, while backward overlaps replica-gradient reduction with ongoing expert computation. MegaFlux extends TensorRT-LLM's CuTeDSL MegaMoE forward kernel and introduces a new backward MoE megakernel. Across 147 configurations per direction on eight NVIDIA B200 GPUs, MegaFlux achieves geometric-mean speedups of $1.45\times$ for forward and $1.28\times$ for backward over the same megakernels with fixed placement, peaking at $2.14\times$ and $2.64\times$. In ablations, pipelining hides $56$--$76$% of replica-weight transfer cost in forward and $91$--$100$% of combined weight-transfer and replica-gradient-reduction cost in backward, yielding up to $13.2$% and $26.7$% additional layer-latency reductions over the same replication plans with these operations executed separately. Integrated into vLLM for DeepSeek-V4-Pro prefill, MegaFlux delivers $1.13$--$1.26\times$ median end-to-end speedups over fixed placement.
展开完整摘要收起摘要↓
Mixture-of-experts (MoE) megakernels fuse expert-parallel communication with expert computation. However, under fixed expert placement, routing skew creates GPU stragglers: overloaded GPUs determine layer latency while others sit idle. Replicating hot experts can shift work to underloaded GPUs, but dynamic replicas introduce additional work: replicas must receive expert weights to execute and, during training, their partial weight gradients must be reduced at the expert owners. We present MegaFlux, which makes expert replication a runtime decision and pipelines the communication induced by replication within persistent MoE execution. An on-device planner jointly selects replica locations and assigns tile-aligned token blocks under a per-GPU replica budget, leaving router outputs unchanged. The forward and backward megakernels realize pipelined expert replication: replicas begin computation as their required weights arrive, while backward overlaps replica-gradient reduction with ongoing expert computation. MegaFlux extends TensorRT-LLM's CuTeDSL MegaMoE forward kernel and introduces a new backward MoE megakernel. Across 147 configurations per direction on eight NVIDIA B200 GPUs, MegaFlux achieves geometric-mean speedups of $1.45\times$ for forward and $1.28\times$ for backward over the same megakernels with fixed placement, peaking at $2.14\times$ and $2.64\times$. In ablations, pipelining hides $56$--$76$% of replica-weight transfer cost in forward and $91$--$100$% of combined weight-transfer and replica-gradient-reduction cost in backward, yielding up to $13.2$% and $26.7$% additional layer-latency reductions over the same replication plans with these operations executed separately. Integrated into vLLM for DeepSeek-V4-Pro prefill, MegaFlux delivers $1.13$--$1.26\times$ median end-to-end speedups over fixed placement.
作者Yang Chen, Yitan Zhang, Michael Witbrock, Shuyue Hu
Inverse Reinforcement Learning (IRL) aims to recover a reward function that explains expert demonstrations. Existing IRL methods typically rely on a bi-level optimization procedure that alternates between reward learning and policy optimization, leading to substantial computational burden and training instability. In this work, we introduce a different route that eliminates policy optimization entirely by leveraging diffusion policies. Our key insight is that a diffusion policy encodes the action-gradient structure of the optimal soft Q function, enabling reward learning to be cast as a sequence of value recovery problems, thereby allowing us to bypass reward-policy loops inherent in prior IRL methods. Specifically, our method proceeds in three stages: (I) recovering the optimal soft Q function via action-gradient matching and estimating the corresponding soft value function (LogSumExp of Q values) in a way inspired by Gumbel regression; (II) calibrating these soft values by inferring a state-dependent offset; (III) extracting the reward by enforcing Bellman consistency. This leads to Loop-Free Inverse Reinforcement Learning (LFIRL), a fully offline algorithm that operates in a simple, loop-free, and sequential manner. LFIRL is simple to implement and significantly improves training efficiency while maintaining strong reward recovery performance. Empirically, across Maze, Franka Kitchen, Adroit Hand Pen, and Push-T benchmarks, LFIRL achieves 2-3x speedup over the fastest baselines, while matching or surpassing state-of-the-art methods in reward recovery quality.
展开完整摘要收起摘要↓
Inverse Reinforcement Learning (IRL) aims to recover a reward function that explains expert demonstrations. Existing IRL methods typically rely on a bi-level optimization procedure that alternates between reward learning and policy optimization, leading to substantial computational burden and training instability. In this work, we introduce a different route that eliminates policy optimization entirely by leveraging diffusion policies. Our key insight is that a diffusion policy encodes the action-gradient structure of the optimal soft Q function, enabling reward learning to be cast as a sequence of value recovery problems, thereby allowing us to bypass reward-policy loops inherent in prior IRL methods. Specifically, our method proceeds in three stages: (I) recovering the optimal soft Q function via action-gradient matching and estimating the corresponding soft value function (LogSumExp of Q values) in a way inspired by Gumbel regression; (II) calibrating these soft values by inferring a state-dependent offset; (III) extracting the reward by enforcing Bellman consistency. This leads to Loop-Free Inverse Reinforcement Learning (LFIRL), a fully offline algorithm that operates in a simple, loop-free, and sequential manner. LFIRL is simple to implement and significantly improves training efficiency while maintaining strong reward recovery performance. Empirically, across Maze, Franka Kitchen, Adroit Hand Pen, and Push-T benchmarks, LFIRL achieves 2-3x speedup over the fastest baselines, while matching or surpassing state-of-the-art methods in reward recovery quality.
作者Shixuan Liu, Tongli Zhou, Junwei Deng, Pingbang Hu, Jiaqi W. Ma
Training data attribution (TDA) estimates the contribution of individual training examples to model outputs. Most scalable TDA methods rely on per-example gradients, whose computation and use at LLM scale pose challenges in efficiency, compatibility, and extensibility. We introduce dattri-LLM, a TDA library that makes gradient-based attribution more practical at scale. For efficiency, dattri-LLM uses compact gradient representations and dynamically routes gradient operations based on a cost model. For compatibility, its capture mechanism collects per-example gradients from existing training loops that call backward(), without requiring changes to the loop or its configuration. This includes distributed training with DDP and FSDP and pipelines built with HuggingFace Transformers, TRL, and OLMo. For extensibility, dattri-LLM exposes reusable gradient operations and training-time callbacks for implementing attribution methods and applications. These interfaces support a variety of attribution methods, including gradient similarity, curvature-based influence, and trajectory-based methods, as well as applications that act on gradients during training, such as online data selection. On the same hardware and workload, dattri-LLM achieves 3.2x the throughput of the fastest competing library on average, scales multiple attribution methods to 110B-parameter models across four H200 GPUs, and offers superior attribution fidelity-cost trade-offs across a range of models with different model families and scales. The source code of dattri-LLM is available at https://github.com/TRAIS-Lab/dattri-llm.
展开完整摘要收起摘要↓
Training data attribution (TDA) estimates the contribution of individual training examples to model outputs. Most scalable TDA methods rely on per-example gradients, whose computation and use at LLM scale pose challenges in efficiency, compatibility, and extensibility. We introduce dattri-LLM, a TDA library that makes gradient-based attribution more practical at scale. For efficiency, dattri-LLM uses compact gradient representations and dynamically routes gradient operations based on a cost model. For compatibility, its capture mechanism collects per-example gradients from existing training loops that call backward(), without requiring changes to the loop or its configuration. This includes distributed training with DDP and FSDP and pipelines built with HuggingFace Transformers, TRL, and OLMo. For extensibility, dattri-LLM exposes reusable gradient operations and training-time callbacks for implementing attribution methods and applications. These interfaces support a variety of attribution methods, including gradient similarity, curvature-based influence, and trajectory-based methods, as well as applications that act on gradients during training, such as online data selection. On the same hardware and workload, dattri-LLM achieves 3.2x the throughput of the fastest competing library on average, scales multiple attribution methods to 110B-parameter models across four H200 GPUs, and offers superior attribution fidelity-cost trade-offs across a range of models with different model families and scales. The source code of dattri-LLM is available at https://github.com/TRAIS-Lab/dattri-llm.
作者Mengyuan Fan, Peizhuang Cong, Zixiao Huang, Si Xu, Tong Qiao, Yanghao Li, Jing Yang, Tong Yang, Quanlu Zhang, Yu Wang
As model sizes continue to scale, distributed training has become inevitable. Automatic parallelization techniques can derive efficient training parallelism strategies at low cost while achieving superior performance. The difficulty of this problem is jointly determined by the complexity of the model and the underlying compute cluster. Meanwhile, mixture-of-experts (MoE) models are increasingly emerging as the dominant architecture and the rapid evolution of accelerator hardware has made cluster heterogeneity commonplace, posing substantial challenges to automatic parallelization. However, existing approaches typically target either MoE architectures or heterogeneous clusters, failing to generalize to scenarios where both challenges coexist. To this end, we present HAPMoE, a heterogeneity-aware automatic parallelism planner for MoE training. HAPMoE builds a lightweight MoE-aware cost model and efficiently searches a six-dimensional parallel space, producing parallel plans directly deployable on Megatron-LM. Experiments show that HAPMoE improves end-to-end training throughput by up to 3.2$\times$ over baselines across heterogeneous clusters. Its non-uniform pipeline partitioning yields an additional up to 78% gains, and its pruning-enhanced dynamic programming algorithm completes the search within 1 minute, demonstrating high efficiency and practical value in complex hardware environments.
展开完整摘要收起摘要↓
As model sizes continue to scale, distributed training has become inevitable. Automatic parallelization techniques can derive efficient training parallelism strategies at low cost while achieving superior performance. The difficulty of this problem is jointly determined by the complexity of the model and the underlying compute cluster. Meanwhile, mixture-of-experts (MoE) models are increasingly emerging as the dominant architecture and the rapid evolution of accelerator hardware has made cluster heterogeneity commonplace, posing substantial challenges to automatic parallelization. However, existing approaches typically target either MoE architectures or heterogeneous clusters, failing to generalize to scenarios where both challenges coexist. To this end, we present HAPMoE, a heterogeneity-aware automatic parallelism planner for MoE training. HAPMoE builds a lightweight MoE-aware cost model and efficiently searches a six-dimensional parallel space, producing parallel plans directly deployable on Megatron-LM. Experiments show that HAPMoE improves end-to-end training throughput by up to 3.2$\times$ over baselines across heterogeneous clusters. Its non-uniform pipeline partitioning yields an additional up to 78% gains, and its pruning-enhanced dynamic programming algorithm completes the search within 1 minute, demonstrating high efficiency and practical value in complex hardware environments.
Mixture-of-Experts (MoE) has increasingly become a mainstream approach for scaling large language models, as it expands model capacity while keeping computation cost nearly constant. Training large-scale MoE models relies on Expert Parallelism (EP), which distributes expert replicas across GPUs and exchanges tokens through all-to-all communication. The efficiency of EP is often constrained by two system bottlenecks: cross-node token transfers are limited by inter-node bandwidth, while skewed expert workloads lead to imbalanced computation across GPUs. Prior work mitigates these bottlenecks based on per-expert workload statistics, but overlooks the fact that experts could share the communication. In this work, we empirically present the observation that many pairs of experts are frequently co-activated by individual tokens. Motivated by this, we present Cobalt, an efficient MoE training framework that leverages expert co-activation to reduce cross-node traffic and workload imbalance. Cobalt adopts a two-stage expert layout planner that adapts expert layout to the evolving expert co-activation and workload conditions. It periodically co-locates frequently co-activated experts on the same node to reduce the cross-node communication, and performs per-step intra-node adjustment to rebalance the workloads. Subsequently, we develop a communication-aware task assignment method that routes tokens to fewer remote nodes based on the current expert layout. Experiments on 32 B200 GPUs show that Cobalt achieves up to 1.53-2.41 times (1.28-1.89 times on average) of speedup compared to existing MoE training frameworks, while reducing cross-node token traffic by 75.74%-99.26%.
展开完整摘要收起摘要↓
Mixture-of-Experts (MoE) has increasingly become a mainstream approach for scaling large language models, as it expands model capacity while keeping computation cost nearly constant. Training large-scale MoE models relies on Expert Parallelism (EP), which distributes expert replicas across GPUs and exchanges tokens through all-to-all communication. The efficiency of EP is often constrained by two system bottlenecks: cross-node token transfers are limited by inter-node bandwidth, while skewed expert workloads lead to imbalanced computation across GPUs. Prior work mitigates these bottlenecks based on per-expert workload statistics, but overlooks the fact that experts could share the communication. In this work, we empirically present the observation that many pairs of experts are frequently co-activated by individual tokens. Motivated by this, we present Cobalt, an efficient MoE training framework that leverages expert co-activation to reduce cross-node traffic and workload imbalance. Cobalt adopts a two-stage expert layout planner that adapts expert layout to the evolving expert co-activation and workload conditions. It periodically co-locates frequently co-activated experts on the same node to reduce the cross-node communication, and performs per-step intra-node adjustment to rebalance the workloads. Subsequently, we develop a communication-aware task assignment method that routes tokens to fewer remote nodes based on the current expert layout. Experiments on 32 B200 GPUs show that Cobalt achieves up to 1.53-2.41 times (1.28-1.89 times on average) of speedup compared to existing MoE training frameworks, while reducing cross-node token traffic by 75.74%-99.26%.
Low-Rank Adaptation (LoRA) is a widely used approach to parameter-efficient fine-tuning (PEFT), yet a performance gap can remain relative to full fine-tuning (FFT). Many LoRA variants improve the initialization or optimization of low-rank factors. At each training step, however, their first-order weight-space directions are constrained by the current parameterization. We characterize the corresponding LoRA-accessible gradient space and show that it coincides with the tangent space induced by the current LoRA parameterization. This characterization yields an orthogonal decomposition of the full weight gradient at the current model parameters. We term the component orthogonal to this space the normal gradient. Based on this decomposition, we propose GDLoRA (Gradient-Decomposed Low-Rank Adaptation). GDLoRA reconstructs the full weight gradient from forward activations and backward signals, extracts its normal component, and directly updates the base weights with this component, while retaining standard AdamW optimization for the LoRA factors. GDLoRA incorporates complementary normal gradients without increasing standard LoRA's optimizer-state memory budget under matched adapter and optimizer configurations. Experiments on natural language understanding, mathematical reasoning, commonsense reasoning, and image classification show that GDLoRA consistently improves over LoRA and narrows the performance gap to FFT. The code is available at https://anonymous.4open.science/r/GDLoRA.
展开完整摘要收起摘要↓
Low-Rank Adaptation (LoRA) is a widely used approach to parameter-efficient fine-tuning (PEFT), yet a performance gap can remain relative to full fine-tuning (FFT). Many LoRA variants improve the initialization or optimization of low-rank factors. At each training step, however, their first-order weight-space directions are constrained by the current parameterization. We characterize the corresponding LoRA-accessible gradient space and show that it coincides with the tangent space induced by the current LoRA parameterization. This characterization yields an orthogonal decomposition of the full weight gradient at the current model parameters. We term the component orthogonal to this space the normal gradient. Based on this decomposition, we propose GDLoRA (Gradient-Decomposed Low-Rank Adaptation). GDLoRA reconstructs the full weight gradient from forward activations and backward signals, extracts its normal component, and directly updates the base weights with this component, while retaining standard AdamW optimization for the LoRA factors. GDLoRA incorporates complementary normal gradients without increasing standard LoRA's optimizer-state memory budget under matched adapter and optimizer configurations. Experiments on natural language understanding, mathematical reasoning, commonsense reasoning, and image classification show that GDLoRA consistently improves over LoRA and narrows the performance gap to FFT. The code is available at https://anonymous.4open.science/r/GDLoRA.
作者Francois Chaubard, Mykel J. Kochenderfer, Chris Ré
Zero-order optimization (ZO) trains without backpropagation, making it relevant to forward-only hardware and non-differentiable loss, but its gradient variance grows with perturbed dimension, inhibiting large-model training. Sharded Optimization Mixture of Assemblies (SOMA) trains LSTM experts independently on $N$ data clusters using simultaneous perturbation stochastic approximation (SPSA), without exchanging gradients, activations or optimizer state. Its separable loss removes cross-expert perturbation noise at the cost of jointly learned representations across domains. Using 80,000 estimated RTX 5090 GPU-hours, we show modest sharding improves training compute efficiency over all tested monolithic ZO controls. At 8.44M parameters and 150 aggregate GPU-hours, SOMA $N=2$ with 64 perturbations reaches 1.76 test nats/byte, versus 2.00--2.11 for monolithic SPSA at 64, 256 or 1,024 perturbations and 2.21 for EGGROLL. On WikiText-103, these frozen checkpoints reach 2.07, 2.25--2.36 and 2.49, respectively. On a fixed separable objective with equal-size blocks, we prove independent losses reduce relative gradient variance to approximately $1/N$ of a shared-loss estimator's. Holding starting weights, data, perturbations and compute fixed, independent rather than summed losses lower SOMA $N=4$ test loss by 0.035 nats/byte after 1,000 updates across three seeds. Larger ensembles offer a separate inference benefit: at similar model size with top-$k$ routing ($k=4$), SOMA $N=256$ achieves 2.36M tokens/s versus 257k for SOMA $N=8$ ($9.19\times$, including routing), at lower test loss (1.68 versus 1.71), albeit using $59.9\times$ as much aggregate training compute. We release all training and evaluation code and checkpoints.
展开完整摘要收起摘要↓
Zero-order optimization (ZO) trains without backpropagation, making it relevant to forward-only hardware and non-differentiable loss, but its gradient variance grows with perturbed dimension, inhibiting large-model training. Sharded Optimization Mixture of Assemblies (SOMA) trains LSTM experts independently on $N$ data clusters using simultaneous perturbation stochastic approximation (SPSA), without exchanging gradients, activations or optimizer state. Its separable loss removes cross-expert perturbation noise at the cost of jointly learned representations across domains. Using 80,000 estimated RTX 5090 GPU-hours, we show modest sharding improves training compute efficiency over all tested monolithic ZO controls. At 8.44M parameters and 150 aggregate GPU-hours, SOMA $N=2$ with 64 perturbations reaches 1.76 test nats/byte, versus 2.00--2.11 for monolithic SPSA at 64, 256 or 1,024 perturbations and 2.21 for EGGROLL. On WikiText-103, these frozen checkpoints reach 2.07, 2.25--2.36 and 2.49, respectively. On a fixed separable objective with equal-size blocks, we prove independent losses reduce relative gradient variance to approximately $1/N$ of a shared-loss estimator's. Holding starting weights, data, perturbations and compute fixed, independent rather than summed losses lower SOMA $N=4$ test loss by 0.035 nats/byte after 1,000 updates across three seeds. Larger ensembles offer a separate inference benefit: at similar model size with top-$k$ routing ($k=4$), SOMA $N=256$ achieves 2.36M tokens/s versus 257k for SOMA $N=8$ ($9.19\times$, including routing), at lower test loss (1.68 versus 1.71), albeit using $59.9\times$ as much aggregate training compute. We release all training and evaluation code and checkpoints.
作者Francois Chaubard, Mykel J. Kochenderfer, Chris Ré
Backpropagation (BP) dominates deep learning but imposes a massive memory tax. For example, training OPT-30B with Adam requires $\approx$ 600GB of GPU memory (assuming batch size 8 and sequence length 2048). Alternatively, zero-order optimization (ZOO) trains in inference-mode (requiring only $\approx$ 60GB for the same model): no stored activations, no gradients, and no optimizer states. However, ZOO convergence has lagged behind BP. In this work, we evaluate two methods to close this gap. First, we show that reallocating training compute budget from many steps to large effective batch sizes with many perturbations (or probes) but fewer steps, allows 1SPSA (Spall, 1992) to outperform zero order methods like MeZO (Malladi et al., 2023) with less training compute. Next, we introduce 1.5-SPSA, adding a single "clean" forward-pass per step to 1SPSA to calculate a cheap diagonal preconditioner in probe-space, which improves convergence rate and convergence by down-weighting high curvature directions. Benchmarking on 6 post-training datasets on both Qwen3 and OPT model families, we show that 1.5-SPSA achieves State-of-the-Art results over previous ZOO solvers with much less optimization steps. For example, we train OPT-13B (for direct comparison to MeZO) and find 1.5-SPSA achieves +3.1% accuracy on SST-2 over both MeZO and BP in only 70 steps vs. MeZO's 100,000 steps. Finally, we combine an 8-bit-packing random generator, triton fused unpack/apply kernels, and distributed parallelism to achieve fast and stable training of models as large as OPT-30B in-place on commodity GPUs (e.g. A100).
展开完整摘要收起摘要↓
Backpropagation (BP) dominates deep learning but imposes a massive memory tax. For example, training OPT-30B with Adam requires $\approx$ 600GB of GPU memory (assuming batch size 8 and sequence length 2048). Alternatively, zero-order optimization (ZOO) trains in inference-mode (requiring only $\approx$ 60GB for the same model): no stored activations, no gradients, and no optimizer states. However, ZOO convergence has lagged behind BP. In this work, we evaluate two methods to close this gap. First, we show that reallocating training compute budget from many steps to large effective batch sizes with many perturbations (or probes) but fewer steps, allows 1SPSA (Spall, 1992) to outperform zero order methods like MeZO (Malladi et al., 2023) with less training compute. Next, we introduce 1.5-SPSA, adding a single "clean" forward-pass per step to 1SPSA to calculate a cheap diagonal preconditioner in probe-space, which improves convergence rate and convergence by down-weighting high curvature directions. Benchmarking on 6 post-training datasets on both Qwen3 and OPT model families, we show that 1.5-SPSA achieves State-of-the-Art results over previous ZOO solvers with much less optimization steps. For example, we train OPT-13B (for direct comparison to MeZO) and find 1.5-SPSA achieves +3.1% accuracy on SST-2 over both MeZO and BP in only 70 steps vs. MeZO's 100,000 steps. Finally, we combine an 8-bit-packing random generator, triton fused unpack/apply kernels, and distributed parallelism to achieve fast and stable training of models as large as OPT-30B in-place on commodity GPUs (e.g. A100).
The pre-training of Large Language Models (LLMs) is increasingly conducted across multiple data centers. As training scales to a larger number of accelerators, the fraction of time spent on computation decreases, while the fraction spent on communication increases. Therefore, frequent synchronization becomes a growing bottleneck. Local update methods reduce this cost by allowing workers to perform several optimizer steps between synchronizations. Most local update methods set the number of local optimizer steps between synchronizations before training and keep this interval fixed throughout the run. However, the best interval can change during the entire train process. If the interval and optimizer are adapted to the current training state, the communication frequency is reduced while maintaining the training performance. In this work, we introduce AutoLoCo, an adaptive training framework to reduce communication in LLM training. It adapts the local interval using scalar training statistics and corrects each outer update. Our method is motivated by two observations: 1) the appropriate local interval varies across training stages, and 2) changing the number of inner steps per interval creates a mismatch with an unchanged outer optimizer, requiring a correction to the outer update. We optimize this mismatch by correction of the outer optimizer for the momentum and the learning rate using the accumulated inner learning rate. Our experiments under communication constraints demonstrate that AutoLoCo reduces communication frequency by 27% relative to DiLoCo while maintaining training performance.
展开完整摘要收起摘要↓
The pre-training of Large Language Models (LLMs) is increasingly conducted across multiple data centers. As training scales to a larger number of accelerators, the fraction of time spent on computation decreases, while the fraction spent on communication increases. Therefore, frequent synchronization becomes a growing bottleneck. Local update methods reduce this cost by allowing workers to perform several optimizer steps between synchronizations. Most local update methods set the number of local optimizer steps between synchronizations before training and keep this interval fixed throughout the run. However, the best interval can change during the entire train process. If the interval and optimizer are adapted to the current training state, the communication frequency is reduced while maintaining the training performance. In this work, we introduce AutoLoCo, an adaptive training framework to reduce communication in LLM training. It adapts the local interval using scalar training statistics and corrects each outer update. Our method is motivated by two observations: 1) the appropriate local interval varies across training stages, and 2) changing the number of inner steps per interval creates a mismatch with an unchanged outer optimizer, requiring a correction to the outer update. We optimize this mismatch by correction of the outer optimizer for the momentum and the learning rate using the accumulated inner learning rate. Our experiments under communication constraints demonstrate that AutoLoCo reduces communication frequency by 27% relative to DiLoCo while maintaining training performance.