ReCoVer: Resilient LLM Pre-Training System via Fault-Tolerant Collective and Versatile Workload
1University of California, Santa Barbara2Argonne National Laboratory
On a 100k-GPU cluster a failure lands about every 18 minutes, and each restart from a checkpoint takes about 10 of them, leaving less than half of the GPU time for training. ReCoVer therefore avoids restarts altogether: it keeps pre-training running through failures, and every iteration still commits the same number of microbatches as a failure-free run, so training stays on the same trajectory.
load checkpoint
restore
restore
no rework
Abstract
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.
One invariant, three layers
ReCoVer holds one rule under any failure schedule. The replicas that are still alive together contribute exactly the number of microbatches, B, that the failure-free run would use in that iteration:
An iteration's gradient is a sum of B microbatch gradients, so it does not matter which replica computes which term. The replacement microbatches come from the survivors' own shards of an effectively endless data stream, which makes each update statistically the same as in the failure-free run. Three decoupled layers keep the rule.
Fault-tolerant collectives
NCCL and standard MPI abort the whole job when a rank dies. ReCoVer adds two primitives on top of ULFM (User-Level Failure Mitigation): a guarded all-reduce and a consensus barrier.
They check that every member of the communicator is alive before reducing. If one is not, they shrink the communicator to the survivors and hand every rank the same view of what failed, so the error is never fatal.
- Detect
MPIX_Comm_agree - Repair
revoke, then MPIX_Comm_shrink - Record
agree on the failure, return early - Reduce
MPI all-reduce on the checked group
Recovery inside the failing iteration
With large global batches an iteration can take minutes, so throwing it away is costly. Before each gradient bucket is all-reduced, ReCoVer snapshots it and tags it with a world epoch that increases on every communicator shrink.
After a failure, only the buckets reduced under the old membership are restored from their snapshots and reduced again among the survivors. The work already done in this iteration is kept, and no gradient mixes two memberships.
Versatile workload
Every surviving replica takes one role per iteration. A major contributes G microbatches, a minor contributes the remainder, and a spare computes like the others but has its gradients zeroed until a failure promotes it. When a failure leaves no spare, some survivors run one extra microbatch in that iteration; the policy then raises G and reassigns the roles.
Before any failure, all 32 replicas are majors and each contributes 8 microbatches.
How one iteration runs
The flowchart puts the three layers together and follows one iteration on a surviving replica. Gray boxes are the usual training step, and each color is one layer.
- Each replica runs its microbatches. Its role decides whether each gradient is kept or zeroed.
- After the last microbatch, every gradient bucket is snapshotted, tagged with the world epoch and reduced through the guarded all-reduce. With no failure, the optimizer steps as usual.
- If a replica failed, the survivors repair the communicator and continue without it. Nothing restarts.
- If a spare is available, it is promoted. The buckets reduced under the old membership are restored from their snapshots and reduced again before the optimizer step.
- If not, this is a policy boundary step. Some survivors run one extra microbatch, the stale buckets are restored on a side stream during its forward pass, and its all-reduce reduces everything again. After the step the policy advances to a larger G and new roles.
Works with 3D parallelism and HSDP
Losing any rank of a model-parallel replica loses a shard of the model, so ReCoVer treats the whole replica as the unit that survives or fails. NCCL stays inside each replica, and only the cross-replica gradient all-reduce goes through the guarded collectives. The same protocol drops into 3D parallelism (data, tensor and pipeline) and into Hybrid Sharded Data Parallel.
In the figure, rank 1 of replica 1 dies. The NCCL error spreads to the other ranks of replica 1, and a NCCL barrier makes sure they all stop, so the replica aborts as a unit. Replica 2 is untouched until its guarded all-reduce sees the loss and shrinks the group to the replicas that are left.
Results
A 7B LLaMA-style model is pre-trained on C4 with 3D parallelism on 512 A100 40GB GPUs: 64 replicas with 128 microbatches each, so B = 8192 per step. Failures are injected during gradient synchronization, the hardest moment to fail because it leaves some gradients partially reduced, which ReCoVer must restore exactly or the loss curve drifts. They arrive at a fixed interval of five iterations, with the exact step inside each interval sampled at random, and 256 of the 512 GPUs are lost in total.
The training trajectory does not change
With half of the GPUs lost over the run, the loss curve stays on the failure-free NCCL run. Across the 259 steps the two differ by 0.02 on average, and there are no spikes when failures hit.
Show the numbers
| Step | Reference | ReCoVer | Difference |
|---|---|---|---|
| 1 | 11.016 | 11.019 | +0.003 |
| 50 | 7.771 | 7.764 | −0.007 |
| 100 | 6.697 | 6.709 | +0.012 |
| 150 | 5.920 | 5.940 | +0.020 |
| 200 | 5.359 | 5.362 | +0.003 |
| 250 | 4.951 | 4.969 | +0.018 |
Throughput drops, climbs, then passes the failure-free run
Before the first failure ReCoVer runs about 3% below the NCCL run: its gradient sync goes through OpenMPI instead of NCCL and is not overlapped with the backward pass. At the first failure throughput drops to about 1,790 tokens per second per GPU, likely because the reshaped cross-replica group has a topology that suits MPI less well. From there it climbs with every failure, and after the last one it runs 8% above the failure-free reference. Each failure step itself takes longer, since it includes the NCCL timeout and the recovery.
Show the averages
| Phase | Reference | ReCoVer | ReCoVer vs reference |
|---|---|---|---|
| Before the first failure (steps 2–6) | 2,016 | 1,965 | −2.5% |
| Between failures 1 and 2 (steps 8–10) | 2,006 | 1,793 | −10.6% |
| After the last failure (steps 165–259) | 2,019 | 2,188 | +8.4% |
| The 32 failure steps | 236–416 |
Why it can beat the failure-free run
The cross-replica all-reduce costs the same in every iteration. The gradient size is fixed by the model, and ReCoVer issues that all-reduce once per iteration, after the last microbatch's backward pass. Compute, on the other hand, grows with the number of microbatches each survivor runs. With G microbatches per replica, s tokens per microbatch, tc compute time per microbatch and Tr for the all-reduce, each GPU processes
This rises with G and approaches the pure compute rate s / tc. As replicas drop from 64 to 32, G goes from 128 to 256, so the same communication cost is spread over twice the compute. Each surviving GPU really does more useful work per second; the rise is not an artifact of the metric.
The NCCL reference follows the same formula with its own fixed Ginit and its faster all-reduce, so ReCoVer passes it once
that is, once G has grown by more than the ratio between the MPI and NCCL all-reduce costs.
Against checkpoint-restart
Checkpoint-restart pays the same restart and rerun cost after every failure, so its effective throughput stays flat. ReCoVer's rises with each failure and reaches 2.23× the baseline after the eighth. Over the whole run it processes 102M more tokens at 234 GPU-hours, 74.9% more than the baseline, and the gap keeps widening until every replica has failed.
This comparison runs on 128 GPUs (16 replicas, 8 microbatches each), which favors the baseline. Restart costs grow with scale: on a production 100k-GPU cluster a restart alone takes about 10 minutes.
Show the numbers
| Failure | GPUs left | Checkpoint-restart | ReCoVer | Ratio |
|---|---|---|---|---|
| #1 | 120 | 157.7 | 170.3 | 1.07× |
| #2 | 112 | 157.7 | 182.9 | 1.15× |
| #3 | 104 | 157.7 | 196.3 | 1.24× |
| #4 | 96 | 157.7 | 210.8 | 1.33× |
| #5 | 88 | 157.7 | 230.8 | 1.46× |
| #6 | 80 | 157.7 | 250.6 | 1.58× |
| #7 | 72 | 157.7 | 278.1 | 1.76× |
| #8 | 64 | 157.7 | 353.3 | 2.23× |
Show the numbers
| GPU-hours | Checkpoint-restart | ReCoVer |
|---|---|---|
| 50 | 39.1M | 41.7M |
| 100 | 58.7M | 75.2M |
| 150 | 92.3M | 114.3M |
| 200 | 117.4M | 172.8M |
| 234 | 136.8M | 239.3M |
What a single failure costs
Both systems pay for the NCCL timeout that detects the failure. After that, checkpoint-restart spends about 100 seconds saving, restarting and reloading, then reruns everything since the last checkpoint, which grows with the checkpoint interval. ReCoVer's recovery takes about 9 seconds and throws nothing away, so it wins even at N = 2, where checkpoints are taken every other step. Its training segment is a little longer because its steps before the failure are slower, as explained above.
Show every component
| N | System | Training | Timeout, exit | Save | Restart init | Load | Rerun | Recovery | Total |
|---|---|---|---|---|---|---|---|---|---|
| 2 | Checkpoint-restart | 0.5 | 132.8 | 25.5 | 71.0 | 4.0 | 15.2 | 249.0 | |
| 2 | ReCoVer | 0.0 | 142.0 | 9.0 | 151.0 | ||||
| 4 | Checkpoint-restart | 7.8 | 132.0 | 24.2 | 74.7 | 4.0 | 21.3 | 264.0 | |
| 4 | ReCoVer | 7.0 | 143.0 | 10.0 | 160.0 | ||||
| 8 | Checkpoint-restart | 20.4 | 131.9 | 24.6 | 69.8 | 4.0 | 34.3 | 285.0 | |
| 8 | ReCoVer | 22.0 | 142.0 | 9.0 | 173.0 | ||||
| 16 | Checkpoint-restart | 46.2 | 132.4 | 24.8 | 68.3 | 4.0 | 58.3 | 334.0 | |
| 16 | ReCoVer | 51.0 | 142.0 | 9.0 | 202.0 | ||||
| 32 | Checkpoint-restart | 93.7 | 132.3 | 25.3 | 71.0 | 4.0 | 99.7 | 426.0 | |
| 32 | ReCoVer | 110.0 | 141.0 | 9.0 | 260.0 | ||||
| 64 | Checkpoint-restart | 198.7 | 131.7 | 26.3 | 71.0 | 4.0 | 205.1 | 636.7 | |
| 64 | ReCoVer | 235.0 | 141.0 | 10.0 | 386.0 |
The paper reports the same results for the HSDP integration on a 1B model.
Citation
@inproceedings{liu2026recover,
title = {ReCoVer: Resilient LLM Pre-Training System via Fault-Tolerant Collective and Versatile Workload},
author = {Liu, Ziyue and Wang, Zhengyang and Zhang, Ruijie and Maurya, Avinash and Zhou, Hui and
Hovland, Paul and Di, Sheng and Cappello, Franck and Nicolae, Bogdan and Zhang, Zheng},
booktitle = {Advances in Neural Information Processing Systems},
year = {2026}
}