A fat tree is the network shape almost every large GPU training cluster uses: hosts hang off leaf switches, leaves connect to a layer of spine switches, and in very large systems spines connect to a third layer of core switches. The name comes from the idea that links get fatter as you climb toward the root, so the tree never narrows into a bottleneck. In modern clusters nobody builds literally fatter cables; instead many identical switches are wired in parallel, a folded Clos network, so the aggregate bandwidth at each level matches the level below.
For a training job this matters in one precise way. Every step ends in collectives, and a collective is only as fast as the slowest path it uses. The fat tree decides how many GPUs can talk to how many others at full speed at the same time, what happens when two flows land on the same cable, and how much of the cluster a broken switch takes with it. This article derives the sizing math from switch radix, shows which collectives are sensitive to oversubscription, explains why hash-based routing collides and what adaptive routing fixes, and ends with how schedulers and NCCL should place jobs on the tree.
From switch radix to cluster size
Start with one switch model with k ports of equal speed, which is called its radix. NVIDIA's Quantum-2 InfiniBand switch, for example, has 64 ports of 400 Gb/s, so k is 64. In a two-tier fat tree each leaf uses half its ports, k/2, for hosts and the other half for uplinks, one to each spine. Each spine has k ports, so it can reach k leaves. That gives k leaves times k/2 host ports: k squared over 2 host ports with k/2 spines and k leaves. With k = 64 that is 2,048 ports from 96 switches.
When you need more, add a third tier. The classic construction by Al-Fares, Loukissas and Vahdat (2008) groups switches into k pods. Each pod has k/2 edge and k/2 aggregation switches, wired as a complete bipartite graph, and (k/2) squared core switches connect the pods. The result supports k cubed over 4 host ports using 5k squared over 4 switches: 65,536 ports from 5,120 switches at k = 64. Every host can reach every other host through many equal-cost paths, (k/2) squared of them between pods, and that path diversity is what the routing layer has to exploit.
In GPU clusters a host port is a NIC port, not a server. An eight-GPU server with one 400 Gb/s NIC per GPU consumes eight ports, usually on eight different leaves when the fabric is rail-aligned (see rail-aligned topology). So 2,048 ports is 256 servers or 2,048 GPUs for one compute fabric, and that number is the practical ceiling of a two-tier design at this radix.
def fat_tree(k, tiers=2, oversub=1.0):
"""Port and switch counts for a fat tree built from radix-k switches.
oversub is down:up bandwidth at the leaf, 1.0 means non-blocking."""
down = round(k * oversub / (1 + oversub)) # leaf ports facing hosts
up = k - down
if tiers == 2:
leaves = k # each spine port reaches one leaf
spines = up
return dict(ports=leaves * down, switches=leaves + spines)
if tiers == 3: # k pods, non-blocking above the leaf
pods, edge = k, k // 2
ports = pods * edge * down
switches = pods * k + (k // 2) ** 2 # edge + aggregation, plus core
return dict(ports=ports, switches=switches)
print(fat_tree(64)) # {'ports': 2048, 'switches': 96}
print(fat_tree(64, oversub=3)) # {'ports': 3072, 'switches': 80}
print(fat_tree(64, tiers=3)) # {'ports': 65536, 'switches': 5120}
Bisection bandwidth and oversubscription
Bisection bandwidth is the capacity across the worst-case cut that splits the hosts into two equal halves. A non-blocking fat tree has full bisection: N ports at link speed B give N/2 times B across any half-and-half split, so every host can send to a partner on the other side at line rate simultaneously. That is the property the topology is built to deliver.
Oversubscription trades it away for cost. A leaf with 48 ports down and 16 up is 3:1 oversubscribed: if all 48 hosts send off-leaf at once, they share 16 uplinks and get a third of line rate each. Traffic that stays inside a leaf is unaffected. Many storage and front-end networks run at 2:1 or 3:1 because their traffic is bursty and mostly not synchronised. Training compute fabrics are usually built non-blocking because a synchronous collective makes every GPU send at the same moment, which is exactly the worst case an oversubscribed tree is designed not to handle.
Which collectives care
Not every collective suffers equally, and knowing which ones do is how you decide whether a cheaper tree is acceptable for your workload.
| Traffic pattern | Where it is used | Cross-leaf share | Sensitivity to oversubscription |
|---|---|---|---|
| Ring all-reduce | Data-parallel gradients | Only ring edges that cross leaves | Low with good placement |
| Tree all-reduce | Large node counts, small messages | Log-depth, few cross links per level | Low to moderate |
| All-to-all | Mixture-of-experts dispatch, some sequence parallelism | Roughly (N - n)/N of all bytes for n GPUs per leaf | High: bounded by bisection |
| Point-to-point | Pipeline stages | One flow per stage boundary | Low unless stages share an uplink |
A ring visits GPUs in an order NCCL picks from the topology. If the ring walks all GPUs on one leaf before hopping to the next, only one edge per leaf crosses the spine layer in each direction, so the uplinks carry a small fraction of the bytes. All-to-all is the opposite: every GPU sends a slice to every other GPU, so almost all bytes cross leaves and the job runs at bisection speed. A mixture-of-experts model that performs well on a non-blocking cluster can lose a large share of its dispatch bandwidth on a 2:1 tree, while a dense model on the same tree barely notices. The ring and tree cost models are worked through in the all-reduce deep dive.
Routing: ECMP collisions and adaptive routing
Full bisection is a property of the wiring. Getting it in practice depends on routing, because a flow has many equal-cost paths and something has to choose one. Ethernet fabrics traditionally use ECMP: each switch hashes header fields such as the source and destination addresses and ports and picks an uplink from the hash. That works for thousands of small web flows. It works badly for training, where a few hundred long-lived RDMA flows each want a full link. Two flows hashed to the same uplink each get half its bandwidth while another uplink sits idle, and in a ring every GPU waits for the slowest flow.
The collision rate is not small. If m flows are hashed independently onto m uplinks, the expected fraction of uplinks left empty is (1 - 1/m) to the power m, which approaches 1/e, about 37 percent. A short simulation makes the point:
import random
def ecmp_trial(flows, uplinks):
load = [0] * uplinks
for _ in range(flows):
load[random.randrange(uplinks)] += 1
return max(load), load.count(0) / uplinks
random.seed(1)
trials = [ecmp_trial(32, 32) for _ in range(10000)]
print(sum(t[0] for t in trials) / len(trials)) # worst uplink carries ~3.5 flows
print(sum(t[1] for t in trials) / len(trials)) # ~36% of uplinks idleA worst link with three or four flows means the collective runs at roughly a third of line rate even though the fabric has capacity to spare. Three remedies exist. First, more flows per connection: NCCL can open several queue pairs per peer (NCCL_IB_QPS_PER_CONNECTION), which gives the hash more independent draws and smooths the load at the cost of more queue-pair state. Second, adaptive routing, where the switch picks the least-loaded port per packet or per flowlet instead of hashing; InfiniBand switches support it, and Ethernet AI fabrics such as Spectrum-X do it with packet spraying and reordering at the receiving NIC. Third, deterministic routing tables computed for the topology: the InfiniBand subnet manager has routing engines such as ftree and updn that spread destinations across spines so a permutation of traffic does not collide. On RoCE fabrics the congestion-control side of this story is covered in RoCE in depth.
Placement: the software half of the tree
Software sees the fat tree through placement. A scheduler that ignores the topology can scatter a 64-node job across every leaf, forcing every ring edge through the spines; a topology-aware scheduler packs it into as few leaves as possible. Slurm does this with TopologyPlugin=topology/tree in slurm.conf and a topology.conf that describes which nodes sit under which switch:
# topology.conf: generated from the cabling plan, never hand-edited
SwitchName=leaf01 Nodes=gpu[001-032]
SwitchName=leaf02 Nodes=gpu[033-064]
SwitchName=leaf03 Nodes=gpu[065-096]
SwitchName=leaf04 Nodes=gpu[097-128]
SwitchName=spine Switches=leaf[01-04]With that file Slurm prefers allocations under the lowest common switch, and the --switches option lets a job ask for at most a given number of leaves and say how long it is willing to wait for that. Inside the job, NCCL discovers the intra-node topology itself, but it cannot see the fabric; on clouds that publish a topology file, NCCL_TOPO_FILE supplies it. The ordering of ranks also matters: launch so that consecutive ranks share a server and consecutive servers share a leaf, which lets NCCL build rings that cross the spine as rarely as possible.
Placement and in-network reduction interact too. With SHARP the aggregation tree lives in the switches, so a job confined to fewer leaves needs fewer switch resources; see SHARP in-network reductions.
Worked example: 1,024 GPUs, non-blocking or 3:1
Suppose you are sizing a 1,024-GPU cluster from 64-port 400 Gb/s switches with one NIC per GPU. A two-tier non-blocking tree supports 2,048 ports, so 1,024 fits with room to grow: 32 leaves with 32 hosts each and 32 uplinks each, and 32 spines. Bisection bandwidth is 512 times 400 Gb/s, about 205 Tb/s.
Now the finance team proposes a 3:1 leaf to save switches and optics. Each leaf carries 48 hosts and 16 uplinks, so 22 leaves (rounded up from 21.3) and 16 spines. For a dense model using ring all-reduce with good placement, each leaf sends one ring edge each way per channel; NCCL typically runs several channels, but they still occupy a handful of the 16 uplinks, so step time barely moves. For a mixture-of-experts model whose all-to-all moves 2 GB per GPU per step, the picture changes: about 95 percent of those bytes cross leaves, each GPU's cross-leaf share is limited to a third of 400 Gb/s, roughly 17 GB/s instead of 50 GB/s, and the all-to-all goes from about 40 ms to about 115 ms. If the step is 500 ms of compute that cannot overlap with dispatch, that is about a 14 percent slowdown on every step for the life of the cluster, which usually costs more than the switches saved. The arithmetic, not the vendor diagram, should drive the choice.
Failure modes
- A single degraded uplink slows the whole job. A cable running at reduced width or with a high symbol-error rate becomes the slowest ring edge. Monitor per-port error counters and link speed, and drain links that flap rather than letting the subnet manager keep rerouting around them.
- Spine loss removes bisection, not connectivity. Losing one of 32 spines removes about 3 percent of cross-leaf capacity; routing converges around it. Losing a leaf disconnects every port under it, which on a rail-aligned fabric means one NIC on each of many servers, so many jobs see a rank failure at once.
- Topology drift. Re-cabling during repairs makes topology.conf wrong. Jobs still run, just slower, and nobody notices. Regenerate the file from LLDP or the subnet manager's view and diff it against the plan.
- Hash polarisation. If every switch tier uses the same hash function and seed, a flow that collided at the leaf collides again at the spine. Use different seeds per tier, or adaptive routing.
- Fragmentation. Small jobs scattered across leaves leave no leaf-contiguous block for the next large job. Reserve whole leaves for large runs and pack small jobs into the remainder.
Trade-offs
A non-blocking fat tree is the safest topology for unknown workloads: every permutation runs at line rate if routing cooperates. It is also the most expensive in switches and, above all, optics, since every tier adds transceivers. Oversubscription is reasonable for inference fleets, storage and front-end networks, and for training clusters that only run dense models with well-placed rings. Rail-only designs drop the spine for cross-rail traffic entirely and rely on NVLink to move data between rails, which is cheaper still but only works when the communication pattern is known. Dragonfly and torus topologies reduce cable count at very large scale, at the price of non-uniform paths and harder routing. For NDR and XDR fabrics the jump in radix is what keeps clusters of a few thousand GPUs at two tiers, which is worth more than raw link speed because every tier you avoid removes latency, optics and a failure domain.
What to do next
- Write down your switch radix and compute the two-tier and three-tier port ceilings with the function above.
- Classify your workload by collective: mostly ring all-reduce, or heavy all-to-all. Estimate the cross-leaf bytes per step.
- Pick the oversubscription ratio from that estimate, not from the budget alone.
- Generate topology.conf from the cabling plan, enable topology/tree, and check that large jobs land on contiguous leaves.
- Run an all-to-all benchmark at full scale and compare it with the bisection number; a large gap means routing collisions, so try more queue pairs or adaptive routing.
- Alert on per-port errors and link-speed downgrades, and re-audit topology after every repair.