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× 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.
Figures & tables
Figure 1 . Training-job recovery. The infrastructure 1 detects a trainer failure and 2 checks hardware operability, then 3 triggers in-place recovery (HOFs only) or spare-based recovery (either failure class).
Figure 2 . HOF recovery using in-memory checkpointing ( e.g., Gemini ( Wang et al., 2023 ) ), with and without spare-based preinitialization.
Figure 3 . Leto ’s architecture, with a state manager for retained state and a shadow trainer for preinitialization.
Figure 4 . RMP organizes model state into protection groups of D HBM data pages, protected by H1 first-tier and H2 second-tier parity pages in host memory.
Figure 5 . Read, write, and read-write accesses to model state ( M ), activation buffers ( A ), and gradients ( G ) during the compute and update phases of rank r ’s iteration t .
Model Updater
Compute (ms)
Update (ms)
Total (ms)
Reference
12566
47
12613
Sync Snapshotting
13424
82
13506 (+7.1%)
Async Snapshotting
13372
82
13454 (+6.7%)
RMU
12572
111
12683 (+0.6%)
Table 1 . Average fault-free iteration time over 30 iterations on OLMoE ( Table 2 ), split into compute and update phases. Reference is the model updater without crash-consistency mechanisms; Sync and Async Snapshotting are the two strawman solutions.
Figure 6 . RMU chunk-level transactions. Bars mark committed chunks; RS(⋅) denotes Reed–Solomon encoding. The undo buffer preserves C2 while its update is in progress.
Figure 7 . GPU memory usage on OLMoE with global batch size 182. The shadow trainer’s prebuilt state causes an OOM in active training after 177 s.
Model
Total params
Active params
Parallelism
Opt. states sharded
DP
TP
PP
EP
Qwen ( Yang et al., 2025 )
4.0B
4.0B
2
2
1
1
✗
OLMoE ( Muennighoff et al., 2025 )
6.9B
1.3B
2
2
1
4
✓
DeepSeek-MoE ( Dai et al., 2024 )
16.4B
2.8B
8
8
1
16
✗
GPT-MoE ( OpenAI et al., 2025 )
20.9B
3.6B
8
8
1
16
✓
Kimi-MoE ( Team et al., 2025 )
49.1B
3.1B
4
4
4
16
✓
Table 2 . Models and parallelism configurations used in our evaluation. DP, TP, PP, and EP denote the degrees of data, tensor, pipeline ( Narayanan et al., 2021 ) , and expert ( Rajbhandari et al., 2022 ) parallelism, respectively. Opt. states sharded indicates whether optimizer states are sharded across data-parallel ranks ( Rajbhandari et al., 2020 ; Zhao et al., 2023 ) .
Figure 8 . Performance with one HOF every 10 minutes. Hatched bars denote spare-based preinitialization (“ + ”); missing bars indicate unsupported workloads.
Figure 9 . ETTR of Leto and Gemini+ on OLMoE and Kimi-MoE with one HOF every 10 minutes, varying the checkpoint interval; vertical lines mark each system’s ETTR-maximizing interval.
Figure 10 . ETTR of Leto and Gemini+ on OLMoE with one failure every 10 minutes, varying the HOF ratio.
Figure 11 . Throughput on OLMoE with failure arrivals from a public trace and a 50% HOF ratio. Shading marks recovery, dotted segments recomputation, and periodic notches checkpointing.
System
Recovery
Init.
Recomp.
ETTR
Gemini (interval=11)
127 s
59 s
68 s
74.1%
Gemini+ (interval=11)
88 s
35 s
53 s
80.5%
Leto (interval=11)
19 s
15 s
5 s
91.0%
Leto (interval=32)
16 s
13 s
3 s
94.2%
Table 3 . Recovery-time breakdown and ETTR on Kimi-MoE for Gemini , Gemini+ , and Leto at the indicated checkpoint intervals. Interval 11 maximizes ETTR for Gemini and Gemini+ , while 32 does so for Leto ( Figure 9 ).
Figure 12 . Training loss of OLMoE over 100 iterations for the fault-free reference, RMP alone, and RMP with RMU, with a HOF injected during the update phase every 10 iterations.
Figure 13 . GPU memory usage of Leto on OLMoE under memory pressure. The gap between the curves is the shadow’s footprint. Levels denote full preinitialization (L3), GPU runtime only (L2), and no GPU state (L1).
Figure 14 . Simulated ETTR of Leto and Gemini+ with a 10-minute MTTF, varying the HOF ratio across model and cluster sizes.
Appendix figures & tables4 assets
Supplementary material from the paper’s appendix.
Appendix
Workload
I ( Gemini , Gemini+ )
W ( MoEvement , MoEvement+ )
10 min
30 min
1 hr
10 min
30 min
1 hr
Qwen
22
37
53
–
–
–
OLMoE
5
9
12
2
4
4
DeepSeek-MoE
15
25
36
8
8
8
GPT-MoE
12
21
29
4
8
8
Kimi-MoE
11
19
27
4
8
8
Appendix
Table 4 . Tuned checkpoint interval I and sparse window W , in iterations, per workload and MTTF. Dashes indicate that MoEvement is not evaluated on the dense workload, Qwen.
Figure 15 . Simulated ETTR at longer MTTFs, using the model and cluster configurations in Figure 14 .
Model
Total params
Active params
Parallelism
Opt. states sharded
DP
TP
PP
EP
DeepSeek-V2 ( DeepSeek-AI et al., 2024 )
236B
21B
8
8
16
16
✓
DeepSeek-V3 ( DeepSeek-AI et al., 2025 )
671B
37B
128
8
16
64
✓
DeepSeek-V4 ( DeepSeek-AI et al., 2026 )
1.6T
49B
1024
8
16
64
✓
Appendix
Table 5 . Models and parallelism configurations used in our large-scale simulations. DP, TP, PP, and EP denote the degrees of data, tensor, pipeline ( Narayanan et al., 2021 ) , and expert ( Rajbhandari et al., 2022 ) parallelism, respectively. Opt. states sharded indicates whether optimizer states are sharded across data-parallel ranks ( Rajbhandari et al., 2020 ; Zhao et al., 2023 ) .
Qwen
OLMoE
DeepSeek-MoE
GPT-MoE
Kimi-MoE
Predicted (s)
3.464
11.939
5.742
5.198
9.246
Measured (s)
3.556
12.002
5.927
5.363
9.097
Error (%)
−2.6
−0.5
−3.1
−3.1
+1.6
Appendix
Table 6 . Simulator accuracy on iteration time. Predicted and measured Titer (s) for every workload in Table 2 , with communication–computation overlap enabled. Measurements are the median of three independent runs from iteration 11 onward.
State-of-the-art large language model (LLM) training takes tens of thousands of graphics processing units (GPUs) for months and encounters failures across the software and hardware stack. Existing fault-tolerance mechanisms either impose non-trivial overhead during failure-free execution or suffer from prolonged recovery latency, particularly under scenarios where a small subset of compute nodes experience permanent failures. %The tradeoff between failure-free overhead and recovery latency forms a space forms a Pareto frontier We present PHOENIX to simultaneously address both optimization objectives. PHOENIX incorporates a fault-tolerance mechanism that restores LLM training via hot-swapping, namely by replacing failed nodes with spare nodes without terminating the complete job. The hot-swapping of PHOENIX is enabled by two ideas: First, it exploits an off-critical-path in-memory checkpointing mechanism for spatial redundancy. Second, it introduces a communicator reconstruction protocol that replaces failed nodes with spare nodes at runtime. PHOENIX efficiently overlaps the in-memory checkpointing with computation, thus introducing zero overhead during error-free execution. Upon permanent node failures, PHOENIX can rebuild memory states with minimal recomputation by leveraging in-memory checkpoints. We evaluate PHOENIX across scales (up to 512 NVIDIA A100 GPUs) and LLMs (up to 65B parameters), and observe zero checkpoint overhead with hot-swapping recovery completing in under 40 seconds. These results show that PHOENIX simultaneously achieves both zero-overhead error-free execution and extremely low recovery cost.
Pre-training large language models on massive GPU clusters has made hardware faults routine rather than rare, driving the need for resilient training systems. Yet existing frameworks either focus on specific parallelism schemes or risk drifting away from a failure-free training trajectory. We propose ReCoVer, a resilient LLM pre-training system that upholds a single invariant: each iteration keeps the number of microbatches constant, ensuring per-iteration gradients remain stochastically equivalent to a failure-free run. The framework is organized as three decoupled protocol layers: (1) Fault-tolerant collectives that isolate faults from propagating across replicas; (2) in-step fine-grained recovery that preserves intra-iteration progress and prevents gradient corruption; (3) versatile-workload policy that dynamically redistributes microbatch quotas across the survivors. The design is parallelism-agnostic, integrating directly with both 3D parallelism and Hybrid Sharded Data Parallel (HSDP) as a drop-in substrate. We evaluate our implementation on end-to-end pre-training tasks for up to 512 GPUs, ReCoVer successfully preserves the training trajectory from a failure-free reference despite of 256 GPUs lost spread across the run. For comparison with checkpoint-and-restart baselines, ReCoVer demonstrates 2.23× higher effective throughput after successive failures. This advantage results in ReCoVer processing 74.9% more tokens at 234 GPU-hours, with the gap widening as the training prolongs.
Ziyue Liu, Zhengyang Wang, Ruijie Zhang +7
University of California at Santa Barbara · Argonne National Laboratory
Large Language Model (LLM) training is frequently interrupted by a heterogeneous spectrum of failures, from common GPU crashes to catastrophic cluster-wide outages. Existing checkpointing systems rely on monolithic, single-tier storage backend, forcing a trade-off between state-saving overhead and recovery speed. We propose TierCheck, a cluster-aware tiered checkpointing system that aligns storage placement with failure heterogeneity. TierCheck adopts a three-tier design that maintains lightweight differential checkpoints in local and peer memory for fast localized recovery, while asynchronously migrating heavyweight base checkpoints to remote persistent storage. It also ensures strict global consistency across tiers without stalling training, and achieves fast cluster-aware checkpoint restoration during recovery. Evaluations on models up to 40 billion parameters show that TierCheck achieves low training overhead, reduces end-to-end checkpointing time to under 10s, and supports high-frequency checkpointing, ultimately striking an optimal balance between low-overhead persistence and fast recovery.
Shujie Han, Feng Jiang, Patrick P. C. Lee +5
Northwestern Polytechnical University · The Chinese University of Hong Kong · National University of Defense Technology