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.

Four layers, one run: who owns which decisionCluster scheduler: Slurm, Kubernetes + Kueue, Rayallocates nodes as a gang; queues, priorities, preemptionRun supervisorowns the run: preflight, restart, exclusion, budgetLauncher: srun + torchrunstarts 8 ranks per node, rendezvous, envTrainer processesstep loop, NCCL collectives, heartbeats, checkpoint save and loadDurable state: checkpoints + run event lognodeslaunchexit, hangThe scheduler knows nodes; only the supervisor knows the run, its step and its failure history.
Layers of training orchestration. The supervisor turns scheduler events into run decisions.

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
done

Two 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 n

The 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.

ShapeWhat you seeDetectorTypical latency
CrashA rank exits non-zero; srun kills the stepProcess exit codeSeconds
HangEvery rank blocked inside a collectiveNCCL watchdog timeout, step heartbeatThe timeout you configured
SlowStep time creeps up; one rank arrives late at every collectivePer-rank step and wait timingsMinutes, if you look
SilentLoss spike, NaN, gradient norm jumpTraining-health checks in the loopOne 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.

The supervisor's control loopALLOCATEN nodes + sparesPREFLIGHTGPU, ECC, NCCL bwLAUNCHresume from ckptRUNNINGwatch heartbeatsFAILEDcrash, hang, slowTRIAGEattribute + exclude nodeGIVE UPbudget spent: page a humanokeventcrash loopretrybad node fails preflight: swap and re-checkEvery transition is written to an append-only event log, which is also the goodput ledger.
Supervisor states. Hardware failures loop back through allocation; software failures exit to 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 componentNaiveTuned
Detect the failure (crash exit vs hang timeout)6 min2 min
Replace the bad node (requeue vs hot spare)8 min1 min
Start containers, rendezvous, build communicators5 min2 min
Load the last committed checkpoint6 min1 min
Warm up: compile, fill pipelines, first steps4 min2 min
Rework: half the checkpoint interval30 min10 min
Total per failure59 min18 min
Checkpoint stall share (C/T)1.7%1.2%
Wall-clock lost33.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

  1. Write down your current restart path and time each stage on a deliberately killed run.
  2. Turn off in-place worker restarts and route every failure through one supervisor with a restart budget.
  3. Add a two-minute preflight with an all-reduce bandwidth floor and an error-counter check.
  4. Set the collective timeout from your longest legitimate collective, add per-step heartbeats, and enable the NCCL flight recorder.
  5. Keep a spare pool and a quarantine pipeline; never return a node without burn-in.
  6. Make preemption a stop request that ends in a committed checkpoint and exit zero.
  7. Compute weekly goodput from the event log, broken down by cause.
Key takeaway: A long training run survives by turning each failure into a short, measured pause: launch through a rendezvous that cannot mix attempts, preflight every node, detect hangs with timeouts and heartbeats, restart from the last committed checkpoint under a supervisor with a budget, spares and quarantine, and track goodput from an event log so you know which delay to cut next.