An Amazon EC2 UltraCluster is not a product you rent by name. It is the way AWS places large numbers of accelerated instances: thousands of them co-located in one Availability Zone and joined by Elastic Fabric Adapter (EFA) networking in what AWS describes as a petabit-scale, non-blocking network, with Amazon FSx for Lustre beside it for fast shared storage. When you launch enough P5 or Trn2 capacity in the right way, you are running in an UltraCluster. AWS has said P5 instances let customers scale to 20,000 H100 GPUs in this arrangement.
For the engineer who has to make a job run well there, the questions are practical: how to get capacity that is actually close together, how to learn where each instance landed, how to order ranks so the heaviest traffic stays local, what a collective really costs across the fabric, and how to survive the failures that come with thousands of GPUs. Instance-level detail, including how to launch with EFA, is in AWS P5 and P4 instances; this article is about the cluster.
What an UltraCluster contains
AWS lists these instance families as deployed in UltraClusters, with the EFA bandwidth it quotes for each. Each instance has a fast internal interconnect, NVLink for the NVIDIA parts and NeuronLink for Trainium, and EFA is the only path between instances.
| Instance | Accelerators | EFA networking per unit |
|---|---|---|
| P4d | 8 x NVIDIA A100 | up to 400 Gbps |
| P5 / P5e / P5en | 8 x H100 / H200 / H200 | up to 3,200 Gbps |
| P6-B200 | 8 x NVIDIA Blackwell B200 | up to 3.2 Tbps, EFAv4 |
| P6e-GB200 UltraServer | up to 72 GPUs in one NVLink domain (GB200 NVL72) | up to 28.8 Tbps total, EFAv4 |
| Trn1 | 16 x Trainium | up to 1,600 Gbps |
| Trn2 / Trn2 UltraServer | 16 x Trainium2 / 64 across four instances | up to 3,200 Gbps / 12.8 Tbps, EFAv3 |
The jump in the right-hand column matters less than the ratio between the internal domain and EFA. On a P5, the eight GPUs share NVLink bandwidth that is many times the instance's EFA bandwidth, so every parallelism plan should keep its most chatty dimension, normally tensor parallelism, inside one instance. UltraServers widen that inner domain to 64 or 72 accelerators, which changes the plan; the Trainium case is worked through in Trainium in depth.
Getting capacity
Getting thousands of GPUs at once is the hard part. There are three common routes. On-Demand Capacity Reservations hold capacity for as long as you pay for them. Capacity Blocks for ML reserve a fixed number of instances for a fixed window that starts in the future; according to the CLI reference, durations go in one-day steps up to 14 days and in seven-day steps up to 182 days, with up to 64 instances per block. Managed services such as SageMaker HyperPod sit on top of either and add cluster management. Cluster placement groups ask EC2 to pack instances close together, and Capacity Blocks are placed in an UltraCluster for you.
# Find a two-week block of 16 p5.48xlarge starting in early November.
aws ec2 describe-capacity-block-offerings \
--instance-type p5.48xlarge --instance-count 16 \
--capacity-duration-hours 336 \
--start-date-range 2026-11-02T00:00:00Z \
--end-date-range 2026-11-23T00:00:00Z
# Buy the offering you chose, then launch into the reservation it creates.
aws ec2 purchase-capacity-block \
--capacity-block-offering-id cb-0123456789example \
--instance-platform Linux/UNIXInstances are launched into a Capacity Block with the capacity-block market type and a target pointing at its reservation ID. Plan for the block to end on time: instances are reclaimed at the end of the window, so the last checkpoint must land on durable storage before then. Check current flags and limits in the AWS CLI reference before scripting, because these APIs have grown new options, such as UltraServer types, over time.
Where did my instances land?
Close is not the same as adjacent. The EC2 instance topology API, DescribeInstanceTopology, returns for each running instance a list of NetworkNodes, ordered from the top layer to the bottom node the instance attaches to; depending on instance type there are three or four layers. Two instances with the same bottom node are physically closest. Use this to order ranks so that neighbouring data-parallel or pipeline ranks share a node.
import boto3
def instance_topology(ids, region, chunk=50):
ec2 = boto3.client("ec2", region_name=region)
out = {}
for i in range(0, len(ids), chunk):
kwargs = {"InstanceIds": ids[i:i + chunk]}
while True:
resp = ec2.describe_instance_topology(**kwargs)
for inst in resp["Instances"]:
out[inst["InstanceId"]] = inst["NetworkNodes"] # top -> bottom
if not resp.get("NextToken"):
break
kwargs["NextToken"] = resp["NextToken"]
return out
def rank_order(topo, private_ip):
"""Sort hosts so instances sharing lower network nodes are adjacent."""
hosts = sorted(topo, key=lambda iid: tuple(topo[iid]))
groups = {}
for iid in hosts:
groups.setdefault(topo[iid][-1], []).append(iid)
print(f"{len(hosts)} instances on {len(groups)} bottom nodes, "
f"largest group {max(len(g) for g in groups.values())}")
return [private_ip[iid] for iid in hosts] # write as the MPI or torchrun hostfileRun this after launch and before the first collective, and fail fast if the spread is much worse than expected, for example if the job straddles two upper-layer nodes when you asked for a small cluster. Slurm users can turn the same data into a topology file so the scheduler packs jobs by node; the general principle is covered in GPU placement.
The path from NCCL to the wire
Between instances, NCCL uses the AWS OFI NCCL plugin to reach libfabric's EFA provider, which runs AWS's Scalable Reliable Datagram (SRD) transport. SRD sprays packets across many network paths and retransmits lost ones in the Nitro hardware; it does not guarantee in-order delivery, leaving ordering to the layers above. Multipath is why EFA tolerates congestion better than single-path transports. For you this means three checks: the EFA device and libfabric version on the AMI match what the plugin expects, FI_PROVIDER=efa selects EFA rather than TCP, and an NCCL debug log at startup shows the plugin and EFA in use. A job that silently falls back to TCP still runs, only many times slower. How collectives decompose into rings and trees is explained in NCCL collectives.
Worked example: an all-reduce on 2,048 GPUs
Consider 256 P5 instances, 2,048 H100s, training an 8-billion-parameter model with data parallelism and BF16 gradients: 16 GB to all-reduce per step. Full optimizer state would not fit in 80 GB per GPU, so it is sharded across ranks, ZeRO stage 1 or FSDP style; the gradient traffic stays about the same. A hierarchical all-reduce first reduce-scatters inside each instance over NVLink, leaving each GPU with a 2 GB shard, then all-reduces each shard across the 256 instances over EFA, then all-gathers inside the instance again. If each GPU gets an even eighth of 3,200 Gbps, that is 400 Gbps, or 50 GB per second.
def ring_allreduce_s(bytes_per_rank, ranks, gbytes_per_s, efficiency=0.7):
return 2 * (ranks - 1) / ranks * bytes_per_rank / (gbytes_per_s * 1e9 * efficiency)
inter = ring_allreduce_s(16e9 / 8, 256, 50) # ~0.11 s across EFA
print(f"inter-node all-reduce: {inter:.3f} s per step")About 0.11 seconds per step, before overlap with the backward pass. If a step computes for 1.5 seconds and the framework overlaps communication well, the network is nearly hidden. If ranks are ordered badly, or one host fell back to TCP, the slowest link sets the pace for the whole ring. That is why rank order and the startup check pay off. The 70 percent efficiency is an assumption; replace it with a measured nccl-tests result from your own allocation.
Checkpoint storage next to the fabric
The fabric is only half of the cluster; the other half is storage that can absorb a checkpoint from every rank at once. A mixed-precision model trained with Adam usually keeps around 16 bytes of state per parameter: BF16 weights and gradients plus FP32 master weights and two optimizer moments. For the 8-billion-parameter example that is about 128 GB per checkpoint; for a 70-billion-parameter model, about 1.1 TB. With sharded checkpointing each rank writes its own slice in parallel, so the write time is the total size divided by the aggregate throughput the file system sustains.
FSx for Lustre throughput scales with the storage you provision and the deployment type you choose, so size it from the checkpoint, not from the dataset. Measure the write time once, then set the checkpoint interval from it and the failure rate. The classic first-order rule, due to Young and refined by Daly, puts the optimum near the square root of twice the write time times the mean time between failures.
import math
def ckpt_bytes(params, bytes_per_param=16):
return params * bytes_per_param
def write_s(total_bytes, fs_gbytes_per_s):
return total_bytes / (fs_gbytes_per_s * 1e9)
def interval_s(write_seconds, job_mtbf_s):
return math.sqrt(2 * write_seconds * job_mtbf_s) # Young / Daly
w = write_s(ckpt_bytes(70e9), fs_gbytes_per_s=20) # ~56 s, measured value preferred
print(f"checkpoint every {interval_s(w, 4 * 3600) / 60:.0f} min") # ~21 min at a 4 h MTBFThe 20 GB per second and the four-hour mean time between failures are placeholders; use your own measurements. Asynchronous checkpointing, which copies state to host memory and writes in the background, shrinks the effective write time and lets you checkpoint more often for the same overhead.
Failure handling at scale
At thousands of GPUs, something fails every few hours: an Xid error, an ECC fault, a degraded NVLink, an EFA device reset or a host that never comes up. A job that assumes otherwise will spend more time restarting than training. The operational pattern is the same everywhere: gate every node before it joins, keep spares, checkpoint often to fast storage and restart automatically.
- Run a short preflight on each instance: GPU health with DCGM diagnostics, an NVLink bandwidth test, and an EFA check such as libfabric's
fi_pingpongbetween neighbours. - Run nccl-tests across the whole allocation and reject outliers whose bus bandwidth is well below the median.
- Hold a few percent of instances as hot spares inside the same reservation, so a replacement lands on the same fabric.
- Checkpoint to FSx for Lustre on a schedule set by measured failure rate and write time, and copy finished checkpoints to S3 asynchronously.
- Automate restart: detect the hang, cordon the bad host, replace it, rebuild the hostfile with the topology order and resume from the last checkpoint.
Watch power and cooling side effects too; synchronous training produces large, simultaneous load swings, discussed in GPU datacenter power.
Failure modes
- Silent TCP fallback. A missing plugin or wrong provider still trains, at a fraction of the speed.
- Random rank order. Ranks assigned by launch order ignore topology and push neighbour traffic through upper layers.
- Block expiry. The Capacity Block window ends on schedule; a checkpoint still in flight is lost.
- Storage as the bottleneck. Thousands of ranks writing checkpoints at once can saturate an undersized file system.
- Cross-AZ assumptions. An UltraCluster lives in one Availability Zone; spanning zones adds latency and data transfer cost.
- Untested spares. A spare that never passed preflight fails at the worst moment.
Trade-offs
Capacity Blocks give short, predictable access without a long commitment, but you must fit the work into the window. Long reservations cost more in total and give flexibility. Self-managed clusters with ParallelCluster or your own tooling give control; HyperPod gives managed resilience at a premium. UltraServers widen the NVLink or NeuronLink domain and favour larger tensor and expert parallel groups, at the price of a larger failure unit. Choose by the shape of the model and the length of the run, not by the size of the headline cluster.
What to do next
- Decide the capacity route: Capacity Block, reservation or managed service, with the window and instance count written down.
- Bake an AMI with matched driver, EFA, libfabric and NCCL plugin versions and test it on two instances first.
- Add the topology script to job launch and build hostfiles from it.
- Gate every node with preflight and nccl-tests; record per-node results.
- Measure checkpoint write time on FSx and set the interval from failure rate.
- Automate restart and keep tested spares in the same reservation.
- Re-run the collective cost estimate with measured bandwidth before scaling up.