A training cluster is designed backwards from the job it must run. The model size and token budget set the compute; compute and a deadline set the GPU count. The parallel layout sets the network, checkpoint size and failure rate set storage and restart tooling, and the GPU count sets the power bill. Teams that start from a hardware quote and work forwards end up with a fabric that is too thin for data-parallel traffic, or a storage tier that stalls every checkpoint.
Capacity planning, meaning how many GPUs to buy and when, is covered in GPU infrastructure planning. This article designs the cluster itself: the scale-up and scale-out domains, fabric tiers, checkpoint storage, the failure-and-restart budget and bring-up. The running example is a 70-billion-parameter dense transformer trained on 15 trillion tokens on H100-class GPUs, and every number is derived.
Start from the compute budget
Training a dense transformer costs about 6 floating-point operations per parameter per token: 2 for the forward pass and 4 for the backward pass. With N parameters and D tokens, total compute is C = 6ND. The GPU's peak throughput is not what you get; model FLOPs utilisation (MFU), the share of peak spent on the model's own maths, is typically 35 to 45 percent for large dense runs. Meta reported 38 to 43 percent for Llama 3 405B in BF16 on H100s. Goodput, the share of wall-clock time spent making progress rather than restarting, multiplies on top.
def size_run(params, tokens, peak_tflops, mfu, gpus, goodput):
flops = 6 * params * tokens
per_gpu = peak_tflops * 1e12 * mfu # useful FLOP/s per GPU
gpu_hours = flops / per_gpu / 3600
days = gpu_hours / gpus / 24 / goodput
return flops, gpu_hours, days
def young_daly(ckpt_stall_s, mtbf_s):
return (2 * ckpt_stall_s * mtbf_s) ** 0.5 # optimal checkpoint interval
def goodput(ckpt_stall_s, interval_s, restart_s, mtbf_s):
lost = ckpt_stall_s / interval_s + (restart_s + interval_s / 2) / mtbf_s
return 1 - lost
flops, gh, days = size_run(70e9, 15e12, 989, 0.40, 4096, 0.95)
print(f"{flops:.2e} FLOPs, {gh/1e6:.2f}M GPU-hours, {days:.0f} days")
# 6.30e+24 FLOPs, 4.42M GPU-hours, 47 daysThe 989 TFLOPS figure is the H100 SXM dense BF16 peak; the sparse figure on spec sheets is double and does not apply to training. The deadline now decides the GPU count. On 4,096 GPUs the run takes about 47 days at 95 percent goodput. On 2,048 it takes over three months, and on 8,192 it takes about three and a half weeks, if the global batch can grow without hurting convergence.
Memory and the parallel layout
Training state for mixed-precision Adam is about 16 bytes per parameter: 2 for BF16 weights, 2 for BF16 gradients, and 12 for the FP32 master copy and two Adam moments. For 70B that is 1.12 TB before activations, against 80 GB per GPU. So the model must be sharded, and the layout determines the traffic each link carries.
| Dimension | Example | Traffic | Where it must run |
|---|---|---|---|
| Tensor parallel (TP) | 8 | All-reduce or reduce-scatter in every layer, forward and backward | Inside the NVLink domain |
| Pipeline parallel (PP) | 4 | Point-to-point activations between stages, per micro-batch | Across nodes, modest bandwidth |
| Data parallel (DP) | 128 | Gradient reduction once per step, overlappable | Across the whole fabric |
TP 8 times PP 4 places a 70B replica on 32 GPUs, about 35 GB of state per GPU before activations. Sharding optimizer state across the data-parallel group frees more for activations. 128 replicas give 4,096 GPUs. The rule that shapes the hardware: tensor parallelism communicates in every layer and must stay inside the scale-up domain, here the 8 GPUs of an NVLink-connected node. Rack-scale domains such as NVIDIA's GB200 NVL72, with 72 GPUs on one NVLink fabric, widen that domain and let larger TP or expert-parallel groups avoid the scale-out network.
The scale-out fabric
The usual H100 design gives every GPU its own 400 Gb/s NIC, eight per node, on InfiniBand or RDMA-capable Ethernet. NIC k on every node connects to the same group of leaf switches, called rail k. Data-parallel collectives pair GPUs with the same local index, so most traffic stays within a rail; rail-aligned topology explains how NCCL exploits that.
Size the fabric from the data-parallel gradient traffic. Each GPU holds 70B / 32, about 2.19 billion parameters, so 4.4 GB of BF16 gradients. A ring all-reduce sends about 2(N-1)/N times the buffer per GPU, about 8.7 GB per step. At an achievable 40 GB/s from a 400 Gb/s link that is about 0.22 seconds. With a 16-million-token global batch, a step is 6 x 70e9 x 16e6 = 6.7e18 FLOPs, about 4.1 seconds of compute per GPU at 40 percent MFU. Communication is about 5 percent of the step and overlaps with the backward pass. Halve the batch or double the GPUs and the ratio doubles, which is why the fabric should be non-blocking rather than oversubscribed.
Switch radix decides the tier count. A two-tier non-blocking fat tree of 64-port switches reaches 2,048 endpoints: 64 leaves with 32 ports down and 32 up. 4,096 GPU ports need a third tier, which a radix-64 fat tree can stretch to 65,536 endpoints. A non-blocking fat tree carries full bandwidth at every level, so 4,096 endpoints means roughly 4,096 links per tier, over 12,000 links in total, each needing two optical transceivers. Optics are a large part of network cost and a major source of link flaps. GPU pod networks covers bring-up of these domains.
Storage: datasets are easy, checkpoints are not
Reading the dataset is rarely the bottleneck for text. 15 trillion tokens stored as 4-byte IDs is 60 TB. Spread over 47 days, that is about 15 MB/s for the whole cluster. Images and video change this, but for language models the storage tier is sized for checkpoints.
A checkpoint holds weights, the FP32 master copy and both Adam moments: about 14 bytes per parameter, 980 GB for 70B. Gradients are not saved. Writing it synchronously in 30 seconds needs about 33 GB/s aggregate write bandwidth. Asynchronous checkpointing copies state to host memory in a few seconds and writes in the background, so the training stall shrinks and the requirement becomes sustained throughput between checkpoints. The checkpointing deep dive covers sharded, crash-safe saves. Keep checkpoint traffic on the front-end or storage network, not the compute rails, so a save does not slow the gradient reduction running beside it.
Failure rates, checkpoint interval and goodput
At this scale hardware fails during every run. The Llama 3 paper reports 466 job interruptions in a 54-day period on 16,384 H100s. 419 were unexpected, and about 78 percent of those were attributed to confirmed or suspected hardware issues, mostly GPUs and their HBM memory. Meta still reached over 90 percent effective training time. Scaled to 4,096 GPUs, the same rate is about two interruptions a day, a mean time between failures (MTBF) of about 12 hours.
The Young/Daly formula gives the checkpoint interval that balances save cost against lost work: interval = sqrt(2 x stall x MTBF). With a 30-second stall and a 12.4-hour MTBF it gives about 27 minutes. With 10 minutes to detect, replace a node and reload, the goodput function above gives:
mtbf = 12.4 * 3600
t = young_daly(30, mtbf) # ~1,640 s, about 27 minutes
print(round(goodput(30, t, 600, mtbf), 3)) # 0.95: 1.8% saving, 3.2% restart and lost workRestart time is the lever. Cutting it from 10 minutes to 3 raises goodput by about one point, roughly 45,000 GPU-hours over this run. That comes from hot spare nodes, health checks that drain bad nodes before the job lands on them, and checkpoint loading from local or nearby storage. Keep 2 to 5 percent of nodes as spares. GPU hardware faults describes the XID errors, ECC events and NVLink failures to watch.
Power and cooling
An H100 SXM draws up to 700 W, and an 8-GPU server with CPUs, NICs and fans is rated at around 10 kW. 512 nodes is therefore about 5 MW for compute alone. Network, storage and cooling overhead, expressed as power usage effectiveness (PUE), typically take the facility to 6 to 7 MW. Air-cooled H100 racks usually hold two to four servers. Rack-scale systems such as the NVL72 draw over 100 kW per rack and require direct liquid cooling. Power is often the constraint that decides the site, and the site then caps the GPU count. Settle it before choosing the GPU generation.
Training also loads the grid unevenly. Thousands of GPUs pause together for a checkpoint or a collective and then resume, swinging the load by megawatts within seconds. Facility and utility teams need that profile early.
Bring-up and acceptance
A cluster is not ready when the cables are in. Bring-up is a sequence of acceptance tests, each gating the next:
- Burn-in per node: GPU stress, HBM memory test, NVLink bandwidth, and thermals under sustained load. Reject outliers.
- Fabric validation: per-link error counters, NCCL all-reduce and all-to-all bandwidth per rail and across the full system, compared with the theoretical bus bandwidth. NCCL collectives explains the bus bandwidth metric.
- Storage: parallel checkpoint write and read at full scale, plus a restore-and-resume test.
- Scheduler: topology-aware placement, so a job's TP groups land inside nodes and its DP groups inside as few spine hops as possible, plus automatic node draining on health-check failure.
- A canary training run of a small model on the whole cluster for a day, tracking MFU, step-time variance and interruptions before the real run starts.
Failure modes
These problems recur across clusters.
- Stragglers. One slow GPU, often thermally throttled or with a degraded link, sets the step time for all 4,096. Track per-rank step time and drain outliers automatically.
- Oversubscribed spines. A 2:1 tapered fabric looks fine in single-rail tests and loses 10 to 20 percent once data-parallel traffic crosses rails at full scale.
- Checkpoint storms. Every rank writes at once and the file system's metadata servers saturate. Use sharded writers with fewer, larger files.
- Silent data corruption. A GPU produces wrong results without an error. Loss spikes or NaNs appear hours later. Run periodic deterministic self-tests and compare replicas' gradient norms.
- Mismatched software. Different driver, NCCL or firmware versions across nodes cause hangs that look like network faults. Pin and verify versions at job start.
Trade-offs
InfiniBand or Ethernet. InfiniBand offers mature congestion control and in-network reduction. RDMA over Ethernet uses commodity tooling and multiple vendors but needs careful congestion tuning. Both run large clusters today; choose on operational skills and supply as much as on benchmarks.
One big cluster or several smaller ones. A single 4,096-GPU fabric runs one big job; four 1,024-GPU islands are cheaper to build and fail independently, but cap any single job's size.
Bigger scale-up domain or faster scale-out. Rack-scale NVLink keeps more parallelism off the network, at the price of liquid cooling and higher rack density.
What to do next
To design a training cluster for your own run:
- Write down N, D, the deadline and an honest MFU; compute C = 6ND and the GPU count with the calculator above.
- Pick the parallel layout so tensor parallelism fits in the scale-up domain, and compute per-GPU state and gradient bytes.
- Size the scale-out fabric from per-step gradient traffic, keep it non-blocking, and count tiers and optics from switch radix.
- Size checkpoint bandwidth from bytes per parameter and target stall, and use asynchronous saves.
- Estimate MTBF from published failure rates, set the interval with Young/Daly, and budget spares and restart time.
- Confirm site power and cooling before choosing hardware.
- Plan burn-in, fabric validation and a full-scale canary run before the first real job.