Ziyue (Alvin) Liu
NeurIPS 2026

ReCoVer: Resilient LLM Pre-Training System via Fault-Tolerant Collective and Versatile Workload

Ziyue Liu1, Zhengyang Wang1, Ruijie Zhang1, Avinash Maurya2, Hui Zhou2, Paul Hovland2, Sheng Di2, Franck Cappello2, Bogdan Nicolae2, Zheng Zhang1

1University of California, Santa Barbara2Argonne National Laboratory

In short

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.

Checkpoint and restart
abort, reload, redo
GPU 1
GPU 2
GPU 3
1
2
3
4 ✕
5
6
sync ✕
exit, restart,
load checkpoint
1
2
3
4
5
6
sync
step
Every replica stops. The job reloads the last checkpoint and computes all 6 microbatches again.
ReCoVer
repair, restore, adjust, continue
GPU 1
GPU 2
GPU 3
1
2
3
4 ✕
5
6
repair,
restore
repair,
restore
dropped from the group
7
8
sync
sync
step
step
no restart,
no rework
The survivors repair the group and each run one extra microbatch: 2 × 3 = 6, the same count as 3 × 2 without the failure.
microbatchextra microbatchwork done againfailure
A GPU fails in the middle of an iteration with three data-parallel replicas and two microbatches each. Numbers are microbatch indices.
2.23×
the effective throughput of checkpoint-restart after eight successive failures
+74.9%
tokens processed at 234 GPU-hours, and the gap keeps growing
256 of 512
GPUs lost in a 7B run with no change to the loss curve

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:

∑r ∈ survivors(t) Cr(t) = Winit · Ginit = B

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.

Bottom layer

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.

  1. DetectMPIX_Comm_agree
  2. Repairrevoke, then MPIX_Comm_shrink
  3. Recordagree on the failure, return early
  4. ReduceMPI all-reduce on the checked group
Middle layer

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.

Top layer

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.

Versatile workload: 32 replicas, 8 microbatches each, B = 256
microbatches committed
32 × 8 = 256

Before any failure, all 32 replicas are majors and each contributes 8 microbatches.

contributes extra microbatch computed, zeroed spare failed promoted spare

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.

Training stepBottom layer: collectivesMiddle layer: gradient recoveryTop layer: workload policy
Computemicrobatch mLastmicrobatch?Accumulate or zeroby roleSnapshot bucket,tag epochGuardedall-reduceFailure?OptimizerstepAdvance policyboundary steps onlyRepaircommunicatorSpareavailable?PromotespareRestore buckets(blocking)Run one extramicrobatchRestore buckets(non-blocking)yesnonoyesyesre-reduceno: policy boundary stepextra microbatch12345Computemicrobatch mLastmicrobatch?Accumulate orzero, by roleSnapshot bucket,tag epochGuardedall-reduceFailure?RepaircommunicatorOptimizerstepAdvance policyboundary steps onlySpareavailable?PromotespareRestore buckets(blocking)Run one extramicrobatchRestore buckets(non-blocking)yesnonoyesyesre-reducenoboundary stepextra microbatch12345
  1. Each replica runs its microbatches. Its role decides whether each gradient is kept or zeroed.
  2. 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.
  3. If a replica failed, the survivors repair the communicator and continue without it. Nothing restarts.
  4. 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.
  5. 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.
Integration

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.

rank 1rank 1rank 2rank 2rank 3rank 3rank 4rank 4Replica 1aborts as a unitReplica 2keeps trainingNCCL: error spreadsNCCLguardedcross-replicaall-reduce

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.

Training loss, 7B model on 512 GPUs
256 GPUs lost to failures injected during gradient synchronization
Failure-free reference (NCCL)ReCoVer, 256 GPUs lost
Show the numbers
StepReferenceReCoVerDifference
111.01611.019+0.003
507.7717.764−0.007
1006.6976.709+0.012
1505.9205.940+0.020
2005.3595.362+0.003
2504.9514.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.

Effective throughput at every step, 512 GPUs
tokens per second per surviving GPU, from step 2 since step 1 includes start-up
Failure-free reference (NCCL)ReCoVerfailure step
Show the averages
PhaseReferenceReCoVerReCoVer vs reference
Before the first failure (steps 2–6)2,0161,965−2.5%
Between failures 1 and 2 (steps 8–10)2,0061,793−10.6%
After the last failure (steps 165–259)2,0192,188+8.4%
The 32 failure steps236–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

tokens/s/GPU = G sG tc + Tr = stc + Tr / G

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

TMPIr / G < TNCCLr / Ginit

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.

Effective throughput after each failure
tokens per second per surviving GPU, 128 GPUs at the start
ReCoVerCheckpoint-restart, flat at 157.7
0100200300400
170.3
353.3
#1#2#3#4#5#6#7#8
failure
Show the numbers
FailureGPUs leftCheckpoint-restartReCoVerRatio
#1120157.7170.31.07×
#2112157.7182.91.15×
#3104157.7196.31.24×
#496157.7210.81.33×
#588157.7230.81.46×
#680157.7250.61.58×
#772157.7278.11.76×
#864157.7353.32.23×
Tokens processed per GPU-hour spent
+102M tokens (+74.9%) at 234 GPU-hours
ReCoVerCheckpoint-restart
Show the numbers
GPU-hoursCheckpoint-restartReCoVer
5039.1M41.7M
10058.7M75.2M
15092.3M114.3M
200117.4M172.8M
234136.8M239.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.

What one failure costs
wall-clock seconds from the last checkpoint until training is back where it was, one failure at step 1.5N
training before the failureNCCL timeout and exitcheckpoint save, restart, reloadrerun lost workReCoVer recovery
N = 2restart
249 s
ReCoVer
151 s
N = 4restart
264 s
ReCoVer
160 s
N = 8restart
285 s
ReCoVer
173 s
N = 16restart
334 s
ReCoVer
202 s
N = 32restart
426 s
ReCoVer
260 s
N = 64restart
637 s
ReCoVer
386 s
0100200300400500600700
seconds
Show every component
NSystemTrainingTimeout, exitSaveRestart initLoadRerunRecoveryTotal
2Checkpoint-restart0.5132.825.571.04.015.2249.0
2ReCoVer0.0142.09.0151.0
4Checkpoint-restart7.8132.024.274.74.021.3264.0
4ReCoVer7.0143.010.0160.0
8Checkpoint-restart20.4131.924.669.84.034.3285.0
8ReCoVer22.0142.09.0173.0
16Checkpoint-restart46.2132.424.868.34.058.3334.0
16ReCoVer51.0142.09.0202.0
32Checkpoint-restart93.7132.325.371.04.099.7426.0
32ReCoVer110.0141.09.0260.0
64Checkpoint-restart198.7131.726.371.04.0205.1636.7
64ReCoVer235.0141.010.0386.0
N is the checkpoint interval in steps. Checkpoint-restart saves every N steps and the failure lands halfway between two saves. 128 GPUs, 8 microbatches per replica.

The paper reports the same results for the HSDP integration on a 1B model.

Citation

BibTeX
@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}
}