A frontier pre-training run is not one job. It is a long sequence of jobs that each die, and an orchestration layer that turns those deaths into short pauses. Meta's Llama 3 paper reports 419 unexpected interruptions during a 54-day snapshot on 16,384 H100 GPUs, roughly one every three hours. At that rate the question is never whether the run fails; it is how many minutes each failure costs, and who notices.
This article is about the control loop around a single training run: launching ranks and rendezvous, proving nodes healthy before trusting them, detecting crashes, hangs and slow nodes, restarting automatically with spares, handling preemption, and measuring the result as goodput. Cluster-level placement, quotas and gang scheduling are covered in the GPU orchestration overview, and what goes into a checkpoint in the checkpointing deep dive. Here we assume both exist and build the supervisor that uses them.
Four layers around one run
Four layers cooperate. The cluster scheduler (Slurm, Kubernetes with a batch queue such as Kueue, or Ray) hands out a gang of nodes. The launcher starts one process per GPU and wires them together. The trainer runs the step loop and writes checkpoints. Between them sits the run supervisor, the piece most teams build last and need first. It is the only component that knows the run as a whole: which step it reached, which nodes have misbehaved, how many restarts it has spent, and whether the last checkpoint is committed.
The scheduler cannot do this job. To Slurm a failed training job is just a non-zero exit code; it does not know that three crashes at the same step point to a data bug rather than hardware. Keep that knowledge in the supervisor and an append-only event log.
Launch and rendezvous
On Slurm the common pattern is one srun task per node, each running torchrun to spawn eight local ranks. torchrun uses a c10d rendezvous: every agent connects to a store on the head node, they agree on world size and rank assignment, and only then do workers start. The --rdzv-id must be unique per attempt group so a stale agent from a dead attempt cannot join the new one.
#!/bin/bash
#SBATCH --job-name=pretrain-7b
#SBATCH --nodes=64
#SBATCH --ntasks-per-node=1
#SBATCH --gpus-per-node=8
#SBATCH --exclusive
#SBATCH --signal=B:USR1@300 # warn the batch shell 5 minutes before the time limit
#SBATCH --requeue
HEAD=$(scontrol show hostnames "$SLURM_JOB_NODELIST" | head -n 1)
STOP_FILE="$RUN_DIR/stop_requested"
trap 'touch "$STOP_FILE"' USR1 # trainer polls this file and saves, then exits 0
export TORCH_NCCL_ASYNC_ERROR_HANDLING=1 # tear down on collective errors
export TORCH_NCCL_DUMP_ON_TIMEOUT=1 # flight-recorder dump on a hang
export TORCH_NCCL_TRACE_BUFFER_SIZE=2000
srun --kill-on-bad-exit=1 torchrun \
--nnodes="$SLURM_NNODES" --nproc-per-node=8 \
--rdzv-backend=c10d --rdzv-endpoint="$HEAD:29500" \
--rdzv-id="$SLURM_JOB_ID-${SLURM_RESTART_COUNT:-0}" \
--max-restarts=0 \
train.py --resume-from latest_committed --stop-file "$STOP_FILE" &
PID=$!
wait $PID # returns early when the trap fires...
while kill -0 $PID 2>/dev/null; do # ...so keep waiting while the trainer saves
wait $PID
doneTwo choices are deliberate. --max-restarts=0 turns off torchrun's in-place restart: at multi-node scale a failed rank usually means a bad node, and restarting workers on the same node just fails again, so the whole attempt exits and the supervisor decides. Second, the time-limit signal goes only to the batch shell (B:), which drops a stop file. Polling a file once per step is boring and portable; forwarding signals through srun and torchrun into eight Python processes per node is neither. A trapped signal makes wait return at once, so the loop waits again; otherwise the script ends and Slurm kills the trainer mid-save.
Preflight: trust no node
Most crashes in the first minutes of an attempt come from nodes that were already bad. A preflight check runs on every allocated node before the trainer starts and refuses nodes that would poison the run. It should take a minute or two, never twenty.
- Inventory: the expected GPU count is visible, driver and CUDA versions match the image, no stray process holds GPU memory.
- Error counters: no uncorrectable ECC errors or pending row remaps, no recent Xid errors in the kernel log. GPU hardware faults explains what each signal means.
- A small compute kernel: a matrix multiply checked against a known answer catches silent corruption and thermal throttling.
- Collective bandwidth: an all-reduce within the node and across a few neighbours, compared with a floor measured on known-good hardware.
- Storage: the checkpoint path is mounted and writable, and the dataset shard index is readable.
import os, time, torch, torch.distributed as dist
def allreduce_busbw_gbps(nbytes=1 << 28, iters=20):
# run under torchrun across the nodes being checked
dist.init_process_group("nccl")
torch.cuda.set_device(int(os.environ["LOCAL_RANK"]))
x = torch.ones(nbytes // 2, dtype=torch.bfloat16, device="cuda")
for _ in range(5):
dist.all_reduce(x) # warm up, build rings
torch.cuda.synchronize()
t0 = time.perf_counter()
for _ in range(iters):
dist.all_reduce(x)
torch.cuda.synchronize()
sec = (time.perf_counter() - t0) / iters
n = dist.get_world_size()
return nbytes * 2 * (n - 1) / n / sec / 1e9 # bus bandwidth, comparable across nThe factor 2(n-1)/n converts algorithm bandwidth into bus bandwidth, which is what NCCL's own tests report and what stays comparable as n changes; the NCCL collectives article derives it. Reject a node whose result falls well below the fleet floor rather than letting one slow link set the pace for thousands of GPUs.
Detecting crashes, hangs, stragglers and silent faults
Failures reach the supervisor in four shapes, and each needs its own detector.
| Shape | What you see | Detector | Typical latency |
|---|---|---|---|
| Crash | A rank exits non-zero; srun kills the step | Process exit code | Seconds |
| Hang | Every rank blocked inside a collective | NCCL watchdog timeout, step heartbeat | The timeout you configured |
| Slow | Step time creeps up; one rank arrives late at every collective | Per-rank step and wait timings | Minutes, if you look |
| Silent | Loss spike, NaN, gradient norm jump | Training-health checks in the loop | One logging interval |
Hangs are the expensive shape. When one GPU drops off the fabric the others wait inside an all-reduce that never completes. The process-group timeout defaults to ten minutes for NCCL, after which the watchdog tears the process down when async error handling is on. Ten minutes of 16,384 idle GPUs is 2,730 GPU-hours per incident. Shorten it with the timeout argument to init_process_group, but never below your longest legitimate collective, including the barrier around a synchronous checkpoint save, or you will kill healthy runs. Pair it with an application heartbeat: each rank writes its step number to the store or a file every step, and the supervisor alarms when the slowest rank stops advancing. The flight-recorder dump enabled above records the last collectives on each rank, which is how you find the one rank that never arrived.
Stragglers do not fail at all. In synchronous data parallelism the step takes as long as the slowest rank, so one throttled GPU slows every GPU. Log per-rank compute time and time spent waiting in collectives; the straggler is the rank with high compute and low wait, while everyone else shows the reverse.
The restart supervisor
The supervisor is a small state machine with a budget. On every failure it attributes blame, excludes suspect nodes, takes replacements from a spare pool and relaunches from the last committed checkpoint. A crash loop, the same failure at the same step on different hardware, is a software or data problem, and retrying it burns money, so it stops and pages a human.
def supervise(run, max_attempts=50, crash_loop=3):
attempts, history = 0, []
while attempts < max_attempts:
nodes = allocate(run.size, exclude=run.quarantined) # spares first, else queue
bad = preflight(nodes)
if bad:
run.quarantine(bad, reason="preflight")
release(bad)
continue # no attempt consumed
ckpt = latest_committed_checkpoint(run) # never a partial save
attempts += 1
log_event(run, "launch", attempt=attempts, step=ckpt.step, nodes=nodes)
outcome = launch_and_watch(nodes, ckpt) # blocks until exit or hang
log_event(run, "exit", **outcome.as_dict())
if outcome.finished:
return "done"
if outcome.stop_requested: # preemption or time limit
continue
suspects = attribute(outcome) # Xid, ECC, missing heartbeat, timeout dump
run.quarantine(suspects, reason=outcome.kind)
history.append((outcome.kind, ckpt.step))
recent = history[-crash_loop:]
if not suspects and len(recent) == crash_loop and len(set(recent)) == 1:
page_human(run, "same failure at the same step on different nodes")
return "crash_loop"
page_human(run, "restart budget exhausted")
return "budget"Quarantined nodes go to a separate pipeline that runs longer burn-in tests before they return to the pool. Never let a node rejoin just because it passed a reboot: intermittent faults are the common case, and a node that failed twice this week will fail again.
Hot spares are the largest single lever. Keeping a few percent of nodes allocated but idle looks wasteful, but replacing a node from a spare takes about a minute, while going back to the queue can take far longer behind other jobs. On Kubernetes, JobSet and Kueue give the gang semantics; the supervisor logic is the same.
Worked example: where the minutes go
Put numbers on it. Take the Llama 3 failure rate: 77,760 minutes and 419 unexpected interruptions give a job-level mean time between failures of about 186 minutes. The table compares a naive setup with a tuned one, using assumed per-step costs. Lost fraction is (stall / interval) + (minutes per failure / mean time between failures), a first-order model that ignores failures during a restart.
| Cost component | Naive | Tuned |
|---|---|---|
| Detect the failure (crash exit vs hang timeout) | 6 min | 2 min |
| Replace the bad node (requeue vs hot spare) | 8 min | 1 min |
| Start containers, rendezvous, build communicators | 5 min | 2 min |
| Load the last committed checkpoint | 6 min | 1 min |
| Warm up: compile, fill pipelines, first steps | 4 min | 2 min |
| Rework: half the checkpoint interval | 30 min | 10 min |
| Total per failure | 59 min | 18 min |
| Checkpoint stall share (C/T) | 1.7% | 1.2% |
| Wall-clock lost | 33.5% | 10.9% |
The naive setup loses about a third of the cluster. No single item fixes it; the restart path is a pipeline of small delays and each one matters. Note where the big wins come from: rework shrinks only if checkpoints are cheap enough to take often, which needs async saves; node replacement only with spares; load time only with a fast restore tier. Training cluster design covers sizing the spare pool and storage for this.
Preemption and planned stops
Planned interruptions are cheaper than failures because they come with notice. Treat them uniformly: time limits, scheduler preemption, node drains for maintenance and rolling driver upgrades all become a stop request. The trainer finishes its current step, writes a checkpoint, waits for the commit, and exits zero. The supervisor sees a clean stop and relaunches without quarantining anything.
On Slurm, --requeue lets the scheduler put a preempted job back in the queue, and SLURM_RESTART_COUNT tells the new attempt how many times that happened. On Kubernetes the notice is SIGTERM followed by a kill after the pod's termination grace period, so the grace period must exceed a full checkpoint save plus commit. Test it on a real run.
Goodput and observability
Measure goodput: the fraction of allocated GPU time that produced training progress you kept. It is computed from the event log, not guessed. Every launch, exit, checkpoint commit and stop request is an event with a timestamp and step. Lost time per incident is the gap from the last committed step to the failure, plus the time from failure to the first new step past that point. Break it down by cause each week; it tells you whether to spend on spares, faster restore, or shorter timeouts.
Alert on rates rather than events: failures per thousand node-hours, minutes to detect, minutes to the first new step, and preflight rejections.
Failure modes
- Restart storm: no budget or backoff, so a bad image relaunches hundreds of times in an hour and spams the scheduler.
- Resuming from an uncommitted checkpoint: the newest directory is half written; always resolve the latest committed marker.
- Timeout shorter than a save: the barrier around a synchronous checkpoint exceeds the collective timeout and the watchdog kills a healthy run while saving.
- Stale rendezvous: a zombie agent from the previous attempt joins the new one because the rendezvous id did not change.
- Zombie processes: a killed attempt leaves a process holding GPU memory; the next attempt fails with out-of-memory on a healthy node. Preflight should catch it.
- Blame on the wrong node: the rank that reports the timeout is usually a victim waiting on the culprit; attribute from the flight recorder, Xid logs and heartbeats.
- Silent divergence after resume: data loader position or RNG state was not restored, so the run is a different run. Compare loss curves across a forced restart.
Trade-offs
Whole-job restart is simple and robust, but every failure pays the full init and load cost. Elastic training, continuing with fewer nodes, avoids waiting for replacements but changes the global batch or the parallel layout mid-run, which complicates reproducibility and needs resharding checkpoints; most large pre-training runs keep a fixed size and use spares instead.
Slurm, Kubernetes and Ray all work; the supervisor logic above is the same on each, and only the allocate and launch calls change.
What to do next
- Write down your current restart path and time each stage on a deliberately killed run.
- Turn off in-place worker restarts and route every failure through one supervisor with a restart budget.
- Add a two-minute preflight with an all-reduce bandwidth floor and an error-counter check.
- Set the collective timeout from your longest legitimate collective, add per-step heartbeats, and enable the NCCL flight recorder.
- Keep a spare pool and a quarantine pipeline; never return a node without burn-in.
- Make preemption a stop request that ends in a committed checkpoint and exit zero.
- Compute weekly goodput from the event log, broken down by cause.