In January 2022 Meta announced the AI Research SuperCluster, RSC, a cluster built for one job: training research models too large for its general fleet. The headline numbers are easy to quote and hard to use. This article treats RSC as a design case study instead. You will see how its parts fit together, from the GPU node through the InfiniBand fabric to the storage tier, and why each choice matters to the software that runs on it. We cover the isolation model, what the Llama 2 team reported when it trained on RSC alongside an Ethernet cluster, and the reliability arithmetic that rules any job spanning thousands of GPUs.
Sources are Meta's launch blog, the Llama 2 paper and a 2024 reliability study; our own arithmetic and assumptions are labelled as such.
What RSC is and why Meta built it
Meta's previous research cluster, from 2017, held 22,000 NVIDIA V100 GPUs and ran about 35,000 training jobs a day. That is a throughput machine: many medium jobs sharing a big pool. RSC was aimed at the other end. Meta framed it around self-supervised learning and transformer models whose size keeps growing, and said it wanted to train models with more than a trillion parameters on data sets as large as an exabyte.
At launch RSC had 760 NVIDIA DGX A100 systems, 6,080 GPUs. Meta's early benchmarks claimed computer vision workflows up to 20 times faster than its older infrastructure, NCCL collectives more than nine times faster, and large NLP training three times faster, so a model with tens of billions of parameters finished in three weeks instead of nine. Phase two was planned to reach 16,000 GPUs by mid-2022, at which point Meta said RSC would deliver nearly 5 exaflops of mixed-precision compute and be the fastest AI supercomputer in the world. Treat both as Meta's own claims made at announcement; the blog describes a plan, and this page does not assume when or how exactly it was completed.
The interesting part is not the GPU count but the balance. A cluster for large jobs must move gradients between nodes as fast as it computes them, and feed data as fast as the GPUs consume it. RSC is three subsystems sized against each other: compute, fabric and storage.
The node: eight GPUs and eight adapters
Each DGX A100 holds eight A100 GPUs joined by NVSwitch, so any GPU reaches any other in the node at full NVLink bandwidth, 600 GB/s per GPU (both directions) on that generation. For the cluster network each system has eight 200 Gb/s HDR InfiniBand adapters, one per GPU, which is where Meta's figure of 1600 Gb/s per system comes from. Separate adapters serve storage.
Software sees a two-level machine. Inside a node, bandwidth is about twelve times higher than between nodes (300 GB/s each way against 25 GB/s per GPU), so the chattiest parallelism goes there: tensor parallelism and the inner part of hierarchical collectives. Across nodes, each GPU talks through its own adapter. One adapter per GPU means NCCL can give every GPU its own path into the fabric instead of funnelling eight GPUs through a shared card. That is the precondition for the NCCL speed-up Meta reported. The all-reduce cost model explains why per-GPU injection bandwidth sets the floor for gradient synchronisation.
The fabric: a non-blocking two-level Clos
Meta describes the fabric as an NVIDIA Quantum InfiniBand two-level Clos with no oversubscription, planned to grow to 16,000 ports in phase two. No oversubscription means each leaf switch has as much bandwidth going up to the spines as coming down from the hosts, so any half of the endpoints can talk to the other half at full rate. That matters because a large training job is placed wherever free nodes exist, and its collectives then cross the spine constantly. With oversubscription, the same job runs at different speeds depending on where the scheduler happened to put it.
Two levels also keep the hop count low: any two hosts are at most leaf, spine, leaf apart. Here is the arithmetic that makes 16,000 ports in two layers plausible. It is our sketch, not Meta's published bill of materials. A 40-port HDR leaf split 20 down and 20 up serves 20 hosts. 16,000 endpoints need 800 leaves. Each leaf sends one uplink to each of 20 spine planes, so each spine needs 800 ports. Fixed 40-port switches cannot do that in two layers; a radix-k fat tree of fixed switches tops out at k squared over two, 800 hosts for k = 40. Director-class modular chassis with hundreds of ports can. The fat-tree radix arithmetic works through the same formula.
The data path: tiered flash and AIRStore
The storage tier at launch was 175 PB of Pure Storage FlashArray, 46 PB of cache in Penguin Computing Altus systems and 10 PB of Pure Storage FlashBlade. Meta built its own service on top, the AI Research Store or AIRStore, which adds a data preparation phase: data sets are preprocessed before training rather than decoded on the fly from raw objects. The target delivery rate was 16 TB/s, with capacity growing toward an exabyte.
A sanity check: at 16,000 GPUs, 16 TB/s is 1 GB/s per GPU. A language model reading tokenised text needs a tiny fraction of that; a video or image model reading compressed frames at high throughput can need a meaningful part of it, and several jobs share the tier. The cache layer exists because training reads the same shards many times across epochs and restarts, and a restart after a failure reloads both data and the latest checkpoint at once. Checkpoint writes are the other large consumer: a model with tens of billions of parameters and its optimiser state is hundreds of gigabytes, and writing it every hour or so must not stall the GPUs.
Isolation by design
Meta wanted to train on data from its own production systems, so RSC was designed around isolation. The blog states that the data path from storage to the GPUs is encrypted end to end, and that RSC is isolated from the internet with no direct inbound or outbound connections. Before data is imported, it goes through a privacy review that confirms it was anonymised or otherwise safeguarded.
For engineers this has practical consequences, which apply to any locked-down cluster:
- Every dependency must come from an internal mirror: Python wheels, containers, model weights and tokenisers. A job that calls out to a public hub at startup simply fails.
- Encryption costs CPU or NIC cycles on the storage path, which is one more reason the data preparation step moves work out of the training loop.
- Logs and traces cannot leave the enclave, so observability must be built inside it.
What Llama 2 training revealed
The Llama 2 paper gives a rare public look at RSC under load. Llama 2 was pretrained on RSC and on Meta's internal production clusters, both with A100 GPUs. RSC used NVIDIA Quantum InfiniBand; the production cluster used RoCE, RDMA over Converged Ethernet, on commodity Ethernet switches. Both connected 200 Gbps endpoints. RSC capped GPUs at 400 W and the production cluster at 350 W. The team concluded that RoCE can scale almost as well as InfiniBand up to 2,000 GPUs.
Read that carefully: it covers one workload at up to 2,000 GPUs with equal endpoint bandwidth, not interchangeable fabrics at 16,000. The power cap gap is a reminder that speed is partly a facility decision: 50 W per GPU across thousands of GPUs is real power and cooling budget.
Reliability at thousands of GPUs
In 2024 Meta researchers (Kokolis and colleagues) published a study of reliability in its research clusters covering 11 months, 4 million jobs and over 150 million A100 GPU hours. Two findings frame operations. Large jobs are far more exposed to failures, because a synchronous job stops when any one of its GPUs, adapters or nodes fails. Yet small jobs are the majority, so infrastructure must serve both. The authors model mean time to failure as a function of job size and propose estimating the effective training time ratio, ETTR: the share of wall-clock time spent making progress.
You can reproduce the shape of that argument with a first-order model. If component failures are independent, a job's failure rate is the sum of the rates of the GPUs it holds, so its MTTF falls in proportion to its size. Checkpointing turns failures into lost work: on average half a checkpoint interval, plus restart time. The code below is our model, not the paper's exact formula, and the failure rate is an illustrative input you should replace with your own fleet's data.
import math
def job_mttf_hours(n_gpus, fails_per_1k_gpu_days):
"""Independent failures: the job fails when any GPU it holds fails."""
return 1 / (n_gpus * fails_per_1k_gpu_days / 1000 / 24)
def ettr(mttf_h, ckpt_write_min, restart_min, interval_min=None):
"""Productive share of wall time under periodic checkpoints (first order)."""
mttf = mttf_h * 60
if interval_min is None: # Young's optimum interval
interval_min = math.sqrt(2 * ckpt_write_min * mttf)
lost = interval_min / 2 + restart_min # expected rework + restart per failure
overhead = ckpt_write_min / interval_min + lost / mttf
return max(0.0, 1 - overhead), interval_min
for n in (2048, 16000):
mttf = job_mttf_hours(n, fails_per_1k_gpu_days=0.5) # assumption, not Meta data
ratio, every = ettr(mttf, ckpt_write_min=5, restart_min=20)
print(f"{n:>6} GPUs: MTTF {mttf:5.1f} h, checkpoint every {every:4.0f} min, ETTR {ratio:.2f}")
Worked example: one job at two scales
Run the model above and the assumed rate of 0.5 failures per 1,000 GPU-days gives a 2,048-GPU job an MTTF of 23.4 hours. With 5-minute checkpoint writes and 20-minute restarts, the best interval is about 119 minutes and ETTR is 0.90. Scale the same job to 16,000 GPUs and MTTF drops to 3.0 hours, the interval to 42 minutes and ETTR to 0.65: a third of an expensive cluster's time goes to checkpoints, rework and restarts. Nothing about the hardware changed; only the job size did.
Now add the communication side. Suppose a data-parallel job synchronises 26 GB of bf16 gradients per step (13 billion parameters). A hierarchical all-reduce first reduces inside each node over NVSwitch, leaving each of the eight GPUs responsible for an eighth, 3.25 GB, which it all-reduces across nodes over its own 25 GB/s adapter. A ring all-reduce moves about twice the buffer per GPU, so the inter-node phase takes roughly 2 x 3.25 / 25 = 0.26 seconds at line rate, or about a third of a second at a realistic 80 percent efficiency, unless it overlaps with the backward pass. Halve the adapters per node and that time doubles. This is the arithmetic behind one adapter per GPU and a non-blocking fabric: they make step time independent of placement.
Two levers follow. To raise ETTR, cut checkpoint write time (sharded, asynchronous checkpoints) and restart time (warm spares, fast health checks). To protect step time, keep collectives on rails and overlap them with compute, as the rail-aligned topology guide explains.
Failure modes
Failures you should expect on any RSC-style cluster:
- One bad GPU or adapter stalls thousands. A synchronous job hangs at the next collective. Without a watchdog that aborts the communicator and restarts from a checkpoint, the hang can burn hours of the whole allocation.
- Placement-dependent speed. If any part of the fabric is oversubscribed, or a cable runs degraded, jobs slow down depending on which nodes they land on. Run collective benchmarks per rack on a schedule.
- Storage stampedes. After a fleet-wide incident, hundreds of jobs restart and read checkpoints at once. Stagger restarts or size the cache tier for it.
- Silent data corruption. A GPU that computes wrong values without crashing can push a loss spike into every replica. Track loss and gradient norms per step and keep known-good checkpoints.
Trade-offs
| Decision | What RSC chose | Gain | Cost |
|---|---|---|---|
| Fabric | InfiniBand, two-level Clos | Mature RDMA, managed routing, low hops | Single vendor, higher price than Ethernet |
| Oversubscription | None | Speed independent of placement | Most switch ports and optics per GPU |
| Adapters | One 200 Gb/s per GPU | Per-GPU path for collectives | Eight cluster cables per node |
| Storage | Tiered flash plus cache, custom AIRStore | High read rate, preprocessing offline | Custom service to build and run |
| Isolation | Encrypted path, no direct internet | Can train on sensitive data | Mirrors and in-enclave tooling required |
| Ownership | On-premises, purpose-built | Full control, predictable cost at scale | Years of lead time, supply-chain risk |
Meta took RSC from a shared document to a working cluster in about eighteen months, through pandemic supply-chain shortages. Renting trades that lead time for higher unit cost and less control, as the OCI GPU SuperCluster walkthrough shows from the other side.
What to do next
- Write down your cluster's three ratios: intra-node to inter-node bandwidth per GPU, adapters per GPU, and leaf uplink to downlink. If the last is below 1, measure job speed by placement.
- Run nccl-tests all-reduce across every pair of racks and keep the results as a baseline; rerun after any fabric change.
- Rerun the ETTR model with your own failure rate, checkpoint time and restart time, for your largest planned job.
- Cut checkpoint write time first (sharded, asynchronous checkpoints), then restart time (spares and pre-flight health checks).
- If the cluster is isolated, mirror every dependency and test a cold start without internet before the first big run.
- Size storage read bandwidth for the restart stampede, not the steady state.
- Read how Meta designs its own accelerator in the MTIA deep dive to see the same balance questions at chip level.