Colossus is the training cluster xAI built in a former appliance factory in Memphis, Tennessee, to train its Grok models. NVIDIA's October 2024 announcement put it at 100,000 Hopper GPUs, said the facility was built in 122 days and that training began 19 days after the first rack was installed, and stated that xAI was doubling it to 200,000 Hopper GPUs. It was the largest publicly described single training cluster at the time.
News coverage focused on speed and controversy. This page does something more useful for an engineer: it reads the published design as a set of choices about how a training job will use the hardware. What fits in a server, what crosses the network, what a rack or an array means to a scheduler, how often something breaks at this scale, and what that does to checkpointing. Where a figure comes from a press report rather than an operator statement, the text says so; xAI has published little itself.
The building blocks
ServeTheHome toured the hall with Supermicro and described the building blocks. The unit is a Supermicro 4U liquid-cooled server holding an eight-GPU H100 HGX board. Eight of those servers and a coolant distribution unit (CDU) make a rack of 64 GPUs. Racks are grouped eight at a time into arrays of 512 GPUs with their networking, described as mini-clusters inside the larger system.
Why does this hierarchy matter to software? Because a training framework maps its parallelism dimensions onto it. Tensor parallelism needs the highest bandwidth and runs inside one server over NVLink, so its degree is at most 8. Pipeline and data parallelism cross servers over Ethernet. If the scheduler places a pipeline's stages in the same rack or array, its point-to-point sends stay below the top tiers of the switch fabric. The same structure appears in DGX H100 system architecture; Colossus is that pattern replicated at enormous scale.
The network: Ethernet at 100,000 GPUs
The choice that made Colossus notable to network engineers was Ethernet. Most earlier clusters of this size used InfiniBand. Colossus used NVIDIA's Spectrum-X platform: Spectrum SN5600 switches (ports up to 800 Gb/s) and BlueField-3 SuperNICs, running RDMA over Converged Ethernet. A reply from ServeTheHome described each server as having one 400 GbE link per GPU plus another for the CPU, which matches the common rail-optimised layout: GPU k in every server connects to the same rail of switches, so collectives that pair GPU k with GPU k elsewhere use one rail.
NVIDIA's release claimed Colossus sustained 95% data throughput with no application latency degradation or packet loss from flow collisions, against about 60% throughput and thousands of flow collisions for standard Ethernet. Those are vendor numbers, but the mechanism is real and worth understanding. Collective communication produces a small number of very large, long-lived flows. Classic ECMP hashing pins each flow to one path; when two elephant flows hash onto the same link, both run at half speed while other links sit idle, and a ring all-reduce runs at the speed of its slowest hop. Spectrum-X sprays packets of a flow across paths (adaptive routing), relies on the SuperNIC to put out-of-order packets back in order, and uses telemetry-driven congestion control to slow senders before buffers overflow. Read NVIDIA Spectrum-X and RoCE in depth for the mechanics.
For the training job, this shows up as predictable all-reduce and all-gather time. A useful back-of-envelope: a ring all-reduce of S bytes across N ranks sends about 2S(N-1)/N bytes per rank. At 400 Gb/s (50 GB/s) per GPU and 95% effective throughput, reducing 2 GB of BF16 gradients takes about 4 GB / 47.5 GB/s, roughly 84 ms per step if nothing overlaps; at 60% it would be about 133 ms. Multiply by tens of thousands of steps and the difference is days of cluster time.
Power and why training load swings
A cluster of 100,000 H100s draws on the order of 150 MW once servers, networking and cooling are counted. Local reporting said the site initially had far less grid capacity than that, that xAI ran on-site natural-gas turbines while utility substations were built (a practice that drew public opposition and permit disputes), and that Tesla Megapack batteries were installed. Treat all of these as reported rather than specified.
The batteries are the part that matters to training engineers. A synchronous training job is a power load that swings in lockstep: every GPU computes, then every GPU waits on a collective, then every GPU computes again. At this scale that is tens of megawatts rising and falling within a fraction of a second, every step, plus a cliff when the job crashes or pauses for a checkpoint. Generators and grid interconnects dislike fast load changes, so a battery between them and the hall can absorb the transients. Software can help too: framework-level power smoothing, staggered restarts and avoiding all-GPUs-idle checkpoint stalls. See datacenter power for GPUs for the electrical side.
Worked example: placing one training job
Put the pieces together with a hypothetical dense model trained with tensor parallelism of 8, pipeline parallelism of 8 and data parallelism across everything else. Each pipeline replica needs 64 GPUs, which is exactly one rack: tensor-parallel groups live inside a server on NVLink, and the eight pipeline stages occupy the rack's eight servers, so activation sends between stages travel one switch hop. An array of 512 GPUs holds eight replicas, and 100,000 GPUs is about 195 arrays, or roughly 1,560 data-parallel replicas.
Gradient synchronisation is the only traffic that must cross arrays. Frameworks handle this hierarchically: reduce-scatter inside the array, all-reduce the shards across arrays on the same rail, then all-gather back. Only a fraction of the bytes crosses the top switch tier, and it can overlap with the backward pass of earlier layers. The layout also turns physical faults into simple accounting. A CDU fault removes one rack, which is one data-parallel replica; with a spare rack in the array, the scheduler can restart the job with the same shape and the same placement rules. Without a spare, the job must either shrink its global batch, which changes training dynamics, or wait.
Real layouts differ: mixture-of-experts models add expert-parallel all-to-all traffic that prefers to stay inside an array, and long-context training adds context parallelism. The method is the same. Write down the bytes each parallelism dimension moves per step, then assign the heaviest to the shortest physical path.
Failure math at this scale
At 100,000 GPUs, hardware failure is a scheduling input, not an exception. The best public calibration is Meta's Llama 3 paper: during a 54-day snapshot of pre-training on up to 16K H100s, the job had 466 interruptions, 419 of them unexpected, with GPU and HBM faults the largest category. That is about 7.8 unexpected interruptions per day. If faults scale with GPU count, a single synchronous job spanning 100,000 GPUs would see about six times as many: roughly 47 per day, or one every 30 minutes.
The Young/Daly rule gives the checkpoint interval that minimises lost work: T = sqrt(2 * C * M), where C is the time to write a checkpoint and M is the mean time between failures.
import math
def plan(gpus, base_gpus=16384, base_unexpected_per_day=419 / 54,
ckpt_write_s=60, restart_s=600):
per_day = base_unexpected_per_day * gpus / base_gpus
mtbf_s = 86400 / per_day
interval = math.sqrt(2 * ckpt_write_s * mtbf_s) # Young/Daly
# fraction of time lost: checkpoint overhead + half an interval redone + restart, per failure
lost = ckpt_write_s / interval + (interval / 2 + restart_s) / mtbf_s
return per_day, mtbf_s / 60, interval / 60, lost
for write_s, restart_s in ((60, 600), (5, 120)): # slow path, then fast path
for g in (16384, 100000, 200000):
d, m, i, l = plan(g, ckpt_write_s=write_s, restart_s=restart_s)
print(f"C={write_s:>2}s R={restart_s:>3}s {g:>7} GPUs: {d:5.1f} failures/day, "
f"MTBF {m:5.1f} min, checkpoint every {i:4.1f} min, ~{l:.0%} time lost")The output is sobering. With a 60-second checkpoint write and a 10-minute restart, the model says a 16,384-GPU job should checkpoint about every 19 minutes and loses about 16% of its time; at 100,000 GPUs the failure interval falls to about 30 minutes, the checkpoint interval to about 8 minutes and the loss to nearly 60%; at 200,000 GPUs the formula exceeds 100%, meaning the job would never make progress. Cut the write to 5 seconds and restart to 2 minutes and the losses fall to about 4%, 14% and 24%. The model is crude: it assumes independent failures that scale linearly with GPU count and that every fault kills the whole job. Meta reported over 90% effective training time on Llama 3 because their real recovery was faster than the slow scenario. The lesson is the direction, not the decimals: at Colossus scale, checkpoint write time and restart time are the levers that decide whether a single synchronous job is viable.
That is why large operators push toward asynchronous, in-memory or peer-replicated checkpoints that cut C to seconds, hot-spare nodes and pre-warmed containers that cut restart time, fast detection of slow or failing GPUs, and sometimes training schemes that tolerate a lost replica instead of stopping. LLM training checkpointing covers the save path itself.
Operating it: failure modes and practices
- Treat each level as a failure domain. A CDU fault can take out a whole rack of 64 GPUs; a leaf switch can isolate many servers. Keep spare capacity per array so a job can be re-placed without crossing more of the fabric.
- Topology-aware placement. Give the scheduler the server, rack and array labels and ask for pipeline stages and expert groups to stay inside an array.
- Burn-in before admission. New racks run GPU, HBM, NVLink and NIC stress tests and an NCCL all-reduce sweep before joining a training job; at 19 days from first rack to training, this has to be automated.
- Straggler detection. One slow GPU or a flapping link slows a synchronous job to its pace. Track per-rank step time and collective time and evict outliers automatically.
- Thermal and coolant telemetry. Liquid-cooled racks fail differently from air-cooled ones: leaks, pump faults and flow drops. Alert on coolant flow and inlet temperature, not just GPU temperature.
Trade-offs
| Choice | Upside | Downside |
|---|---|---|
| Ethernet (Spectrum-X) over InfiniBand | Familiar operations, multi-vendor optics, scale-out tooling | Needs adaptive routing and tuned congestion control to match IB |
| Direct liquid cooling | Higher rack density, less hall air handling | Plumbing failure modes, CDU per rack |
| Build in an existing factory | Months instead of years | Power and permitting become the critical path |
| On-site generation and batteries | Run before the grid is ready, absorb transients | Emissions, permits, community opposition |
| One giant synchronous job | Fastest time to a trained model | Failure rate and checkpoint overhead scale with size |
Lessons for smaller clusters
Few teams will build anything like Colossus, but the lessons scale down. A 512-GPU array is itself a respectable cluster, and its design rules apply at 64 or 256 GPUs: keep tensor parallelism inside NVLink, give each GPU its own network rail, make the scheduler aware of physical domains, budget for failures from day one, and measure collective throughput rather than trusting port speeds. For the sizing method, see LLM training cluster design.
What to do next
- Draw your own cluster as server, rack, array and site, and label each level's failure domain and bandwidth.
- Map tensor, pipeline, data and expert parallelism onto those levels; keep tensor parallel within NVLink.
- Run an NCCL all-reduce benchmark across the full job size and compare bus bandwidth with the per-GPU link rate.
- Estimate your failure rate from your logs (or scale from Llama 3's 419 unexpected interruptions in 54 days on 16K GPUs) and compute a Young/Daly checkpoint interval.
- Measure checkpoint write and restart time; if either is minutes, make it asynchronous or add hot spares.
- Add per-rank step-time and collective-time dashboards so stragglers are found in minutes.