Training a big AI model takes thousands of chips working together for weeks. Chips in the same cabinet, called a rack, are joined by very fast, short links (Scale-up fabrics). But one rack is not enough. So the racks are tied together by a second network that spans the whole building: the .
You can’t run a cable from every chip to every other chip. Ten thousand chips would need fifty million cables! Instead, each chip plugs into a , a box that passes messages along. The switches plug into more switches, in layers.
This chapter asks three questions. How should the switches be arranged? How do chips send data quickly? And what happens when too much traffic heads the same way at once?
A training job spread over thousands of accelerators has to exchange data every step. Within a server or rack, accelerators use dedicated scale-up links (Scale-up fabrics). Between servers, they use a : each accelerator has a with a cable to a , and switches connect to other switches in layers.
Large AI operators keep this “back-end” network separate from the ordinary data-center network that handles storage and management, so the two can evolve on their own schedules.14 A typical AI server today has eight accelerators and eight NICs, one per accelerator, each at 400 gigabits per second (Gb/s).1415 The traffic is unusual by data-center standards: a few very large, periodic, bursty flows per server rather than millions of small ones.15
Four things decide how well the network serves training:
- Topology: how many switches, in how many layers, wired how. This sets how much traffic can cross the network at once and how many switches a message passes through.
- Transport: how data moves between machines. AI clusters use , which lets NICs copy data directly between memories, over either or Ethernet.
- Congestion control: how senders slow down when a switch gets crowded.
- Load balancing: how traffic is spread over the many possible paths.
Why the traffic looks the way it does, and which parts of a model’s communication land on this network, is the subject of Why the network looks this way.
The is the NIC-attached fabric that joins scale-up domains (servers or racks; see Scale-up fabrics) into one cluster. In current AI deployments it is a dedicated back-end network, physically separate from the front-end network that carries storage, checkpointing and management traffic.14 Each accelerator gets its own NIC, typically 400 Gb/s, so an eight-GPU server injects 3.2 Tb/s.1415
What makes it different from a general-purpose data-center network is the traffic. Cloud workloads produce millions of small flows that statistically even out. Training produces a few long, line-rate, synchronized flows per NIC, and every collective waits for its slowest member.15 That shifts the design questions:
- Topology: folded Clos (fat tree) built from identical switch chips, with radix and tier count setting scale, hop count and optics; how much oversubscription to allow and where; and whether to wire by rail.
- Transport: RDMA verbs over InfiniBand, over Ethernet as RoCEv2, or over the Ultra Ethernet transport; lossless links (credits or PFC) versus lossy links with smarter recovery.
- Congestion control and load balancing: ECN- and delay-based rate control, and what to do when per-flow hashing meets a handful of elephant flows.
The chapter takes these in order, then sizes real fabrics in “By the numbers” and works the math in “Under the hood.” The mapping from parallelism strategy to traffic is in Why the network looks this way; the copper and optical links themselves are in Optics.
3 tiers of 64-port switches: 783 switches, 30,032 cables, at most 5 switch hops. Cables grow with n, not n².
One switch is not enough
A switch has a fixed number of sockets for cables. A big one has 64. Plug 32 servers into it and you still have 32 sockets free to connect to other switches. That is the whole trick.
Two layers
Put a row of “lower” switches next to the servers. Above them, put a row of “upper” switches. Give every lower switch a cable to every upper switch. A message then passes three switches at most. With 64-socket switches, this connects about two thousand chips.
Three layers
For more chips, add a third layer on top. Now the same switches connect about sixty-five thousand chips. Engineers call this shape a . A real tree has one trunk, which would be a traffic jam. In a fat tree, the “branches” get thicker toward the top, with more and more cables. So the top never becomes the slow spot.2
Fewer cables going up
Cables and switches cost money. So many networks give each lower switch fewer cables going up than coming in. Say 48 servers share 16 cables going up. If all 48 talk to faraway servers at once, each gets only a third of its speed. Ordinary datacenters have long done this.3 AI clusters usually don’t, because all their chips really do talk at once.14
Rails
An AI server usually has eight chips, each with its own network cable. A network wires them in a clever way. Cable 1 from every server goes to switch 1, cable 2 to switch 2, and so on. Chip 3 in one server often talks to chip 3 in another. With rails, those two are only one switch apart.15
Radix and reach
A switch’s is its number of ports. One switch with radix can connect at most devices. To connect more, you build a network of switches, and the radix decides how many layers (tiers) you need.
Leaf and spine
The standard two-tier design is the network. Leaf switches (often called top-of-rack switches) connect to servers with half their ports and to spine switches with the other half. Every leaf has one link to every spine. With radix :
- each spine has ports, one per leaf, so there can be leaves;
- each leaf has ports down, so the network connects NIC ports;
- any path is leaf → spine → leaf: at most three switch hops;
- between any two leaves there are equally short paths, one through each spine.
With 64-port switches: 64 leaves, 32 spines, 2,048 NIC ports.
Where the design comes from: Clos and fat trees
In 1953 Charles Clos, working on telephone exchanges, showed how to build a big switch from three stages of small ones, using fewer crosspoints than one giant switch. His example connected 36 inputs to 36 outputs with 1,188 crosspoints instead of 1,296.1 A “folded” in half, with the input and output stages merged into one row of leaves, is exactly the leaf–spine design.
In 1985 Charles Leiserson proposed the for supercomputers: processors at the leaves, switches at the inner nodes, and more wires on each link the closer it is to the root, so bandwidth grows going up.2 In 2008, Al-Fares, Loukissas and Vahdat showed how to build one from identical, cheap Ethernet switches. Their three-tier -ary fat tree has “pods,” each with edge (leaf) and aggregation switches, plus core switches on top, and supports hosts. With 48-port switches that is 27,648 hosts.3 Google’s datacenters went the same way: five generations of Clos fabrics built from merchant switch chips, growing 100-fold in ten years to over 1 petabit per second of capacity across the network.4
Bisection bandwidth
How much traffic can a network carry? The usual yardstick is : split the servers into two equal halves in the worst possible way, and add up the capacity of the links that cross between them. A network with full bisection bandwidth lets every server in one half send to a partner in the other half at full speed, all at once. For 2,048 NICs at 400 Gb/s, that is terabits per second (Tb/s).
Oversubscription
Uplinks and the switches above them are expensive, so designers often give each leaf fewer ports up than down. The ratio is the . A 64-port leaf with 48 ports down and 16 up is 3:1: when every server sends across the network at once, each can get at most a third of its line rate. Traffic that stays under one leaf is unaffected. Al-Fares et al. reported that general-purpose datacenters of the time were typically oversubscribed 2.5:1 to 8:1 to save cost.3 AI back-end networks are usually built non-blocking (1:1) within a large block of GPUs and oversubscribed only above it, where little traffic is expected.1415
Rail-optimized wiring
An eight-GPU server has eight NICs. In a network, NIC 0 of every server connects to one leaf, NIC 1 to another, and so on; each set is a “rail.”17 Training software arranges most of its server-to-server traffic between GPUs with the same position (rank) in different servers, so rails make that traffic cross just one switch.16 The effect on scale is large. With all eight NICs on one leaf, a 64-port leaf with 32 ports down serves 4 servers. With rails, a group of eight leaves serves 32 servers, and the 32 GPUs on each rail are one hop apart: eight times as many servers as before. Alibaba reports exactly this 8× gain in its design.15 Traffic between different ranks can go up through the spine; Alibaba instead first crosses to the right GPU over the server’s internal scale-up link.15
Folded Clos, counted
Clos’s 1953 construction has an ingress stage of switches with inputs each, middle switches, and an egress stage mirroring the ingress. Each ingress switch has one link to each middle switch. The worst case for a new connection is that of the ingress switch’s middle links and of the egress switch’s are busy, on disjoint middle switches; so middle switches guarantee a free path (his 36-port example: , ).1 Data-center fabrics fold the network so the ingress and egress switches are the same leaves, and settle for (one uplink per downlink). That makes the fabric rearrangeably non-blocking: for any permutation some set of paths saturates every host link, but finding it is a routing problem, not a guarantee.3
For radix and a non-oversubscribed fabric:
| 2 tiers (leaf–spine) | 3 tiers (-ary fat tree) | |
|---|---|---|
| Endpoints | ||
| Switches | ||
| Switches per endpoint | ||
| Worst-case switch hops | 3 | 5 |
| Equal-cost paths, far endpoints |
The three-tier row is Al-Fares et al.’s construction: pods of edge and aggregation switches, cores, hosts; for , 27,648 hosts and 576 equal-cost paths between hosts in different pods.3 Every switch is the same part, which is the point: Leiserson’s fat tree needed fatter links near the root, and a folded Clos gets the same effect by multiplying identical thin ones.2 Leiserson also proved the shape is no accident: for a given volume of hardware, a fat tree can simulate any other network built from the same hardware with at most a polylogarithmic slowdown.2
Radix is the master variable
Because reach goes as or , a 2× radix increase buys 4× or 8× more endpoints at the same tier count, or removes a tier at the same size. Removing a tier saves two switch traversals on the worst path, 40% of the switches per endpoint ( versus ) and half the switch-to-switch links. That is why the per-chip SerDes budget matters so much: a 51.2 Tb/s chip built from 512 × 100G SerDes can run them as 64 × 800G ports, or group them more narrowly as 128 × 400G.19 At 400G per GPU, the 128-port configuration reaches 8,192 GPUs in two tiers instead of 2,048. Alibaba’s HPN exploits the same high radix, running each tier-1 chip with 128 × 200G ports down toward GPUs, to put 1,024 GPUs under a single layer of rail switches and about 15,000 GPUs in a two-tier pod, a size that “typically” needs three tiers.15
Bisection bandwidth and where to oversubscribe
is the minimum, over all balanced cuts, of the capacity crossing the cut. For a folded Clos with endpoints at rate and leaf ratio (downlinks to uplinks, all same speed), cuts between leaves are bounded by the leaf uplinks: . Al-Fares et al. define as the ratio of worst-case achievable host bandwidth to this bisection; 5:1 means only 20% of host bandwidth is available for some patterns.3
Where you oversubscribe matters as much as how much. Production AI fabrics keep the lower tiers at or near 1:1 and push the ratio to the top, then make the job scheduler respect it. Meta’s “AI Zones” are non-blocking two-stage Clos networks, with an oversubscribed aggregation layer joining zones, and a scheduler that places jobs to minimize cross-zone traffic.14 Alibaba runs tier 1 at 1.067:1 (128 × 200G down, 60 × 400G up per switch) and the aggregation-to-core layer at 15:1, assigning only pipeline-parallel traffic across pods.15
Rails, PXN and rail-only
A rail is the set of GPUs with the same local rank across scale-up domains, and a network puts each rail under its own leaves.16 Two things change:
- Locality. A rail group of eight leaves with downlinks each serves servers, and all same-rank traffic among them is one hop. Without rails the same eight leaves serve servers too, but only per leaf, so the one-hop set is 8× smaller.15
- Cross-rail traffic. GPU 1 in server A talking to GPU 2 in server B either climbs to the spine, or first moves over NVLink (or another scale-up link) to GPU 2 in server A and then rides rail 2.15 NCCL calls the second option PXN; NVIDIA measured all-to-all running more than twice as fast with it.17
Taken to its limit, this argues for removing the spine entirely. Wang et al. found that LLM training traffic is sparse and concentrated within rails, and that a “rail-only” network (a separate Clos per rail, cross-rail traffic forwarded through the scale-up domain) matches a rail-optimized network’s performance for models without mixture-of-experts layers, at 38–77% lower network cost than a full-bisection any-to-any Clos. For mixture-of-experts models, whose all-to-all must be forwarded, they estimate an 8.2–11.2% completion-time overhead on that traffic.16
A → B (another pod): 5 switch hops, 4 equal-cost paths (showing 1).
Without rails each server sits on its own leaf, so any server-to-server traffic goes leaf → spine → leaf: three hops.
Skipping the middleman
In a normal computer, the main processor handles every network message. It packs the data, addresses it and unpacks it at the other end. At AI speeds, that would keep it busy doing nothing but paperwork. So AI clusters use . The network card copies data straight out of one machine’s memory and into the other’s, with no help from the main processor.8
Two kinds of network
Data travels in small chunks called packets. When a switch gets too full, something has to give. , a network built for supercomputers, never lets a switch get too full. A switch only sends to the next one after it hears “I have room.”10
Ethernet is the network almost every other datacenter uses. Normally, a full Ethernet switch just throws packets away, and the sender tries again later. That is fine for web pages but bad for AI. So AI Ethernet adds a “pause” signal: a switch that is filling up tells the one feeding it to wait.7
Both kinds are used to train the biggest models. A newer version, Ultra Ethernet, aims to work without the stop sign at all.12
What RDMA does
With ordinary networking (TCP/IP), the operating system’s code on the CPU builds and parses every packet, copies data between buffers, and handles acknowledgments and retransmissions. Microsoft measured that sending at 40 Gb/s over TCP used 6% of a 32-core server’s CPU, and receiving used 12%.7 moves that work into the : applications register memory buffers with the NIC, and the NICs transfer data between registered buffers directly, bypassing the host networking stack.8 In Microsoft’s production measurements, RDMA’s 99th-percentile latency was 90 microseconds (µs) against 700 µs for TCP on the same network.7 For AI, NICs can also read and write GPU memory directly, so GPU-to-GPU traffic skips host memory entirely.14
The catch is that the NIC’s transport is simple. It was designed assuming the network almost never drops packets.8 So the question becomes: how does each network technology avoid drops?
InfiniBand: credits
was designed around RDMA. Its link layer uses : the receiving end of each link tells the sender how much buffer space it has, and the sender never sends more than that.8 Packets are not dropped for lack of buffer space, and the link delivers them in order, so the transport layer stays simple.10 A central program, the subnet manager, configures the network: it computes routes and programs them into the switches’ forwarding tables, and reconfigures them when a link fails or appears.1011
Ethernet: RoCE and pause frames
(RDMA over Converged Ethernet) runs the same RDMA transport over Ethernet. Version 2, the one used at scale, wraps it in ordinary IP and UDP headers, so it can be routed through a standard data-center network; the destination UDP port is always 4791 and the source port is chosen per connection so switches can spread connections over paths.7
Ethernet switches normally drop packets when a buffer is full. To make Ethernet lossless for RDMA, RoCE networks turn on : when a switch’s buffer for one traffic class crosses a threshold, it sends a pause frame to the port feeding it, which stops sending that class until the buffer drains.7 PFC works, but it pauses everything in that class on that link, including traffic that isn’t headed for the congested spot. This is called , and pauses can spread backward hop by hop across the network and, in rare cases, deadlock it.87
InfiniBand or Ethernet?
Both carry the largest training jobs. InfiniBand offers a lossless fabric designed for RDMA from the start, with centrally managed routing. Ethernet offers many vendors, standard IP routing and tooling, and the same switch chips used across the datacenter, at the cost of retrofitting RDMA’s assumptions onto it. Meta, for example, runs its training clusters on RoCE over Ethernet and built a 24,000-GPU cluster that way.14
The Consortium’s specification, first published in June 2025, takes a different route: a new transport designed for Ethernet that tolerates lost and reordered packets, so it doesn’t need PFC at all. It states that it sets out to address RoCEv2’s shortcomings and scales to millions of endpoints.1213
RDMA’s assumption
NICs implement the transport (segmentation, ordering, acknowledgment, retransmission) in hardware and DMA directly to and from registered buffers, bypassing the host stack.8 NIC resources are tight, so classic implementations keep little per-packet state: in-order delivery is expected and loss recovery is go-back-N, which is cheap only when loss is rare.89 Everything else in this section follows from that assumption.
InfiniBand’s link layer
satisfies the assumption at the link: the receiver on each link advertises credits, and the sender never sends beyond them.89 Since the link is lossless and delivers in order within a virtual lane, the transport can stay simple: it was not designed to recover efficiently from loss and uses go-back-N.109 Routing is computed centrally by a subnet manager that programs paths into every switch and recomputes them on link changes.1011 Current InfiniBand switches add adaptive routing and in-network aggregation and reduction (SHARP), so part of a collective’s arithmetic happens in the switches.20
RoCEv2 and PFC, in detail
carries the IB transport in Ethernet/IPv4/UDP. The destination port is 4791; the source port is random per queue pair, and switches hash the five-tuple, so one QP stays on one path and different QPs can take different paths.7 Losslessness comes from : when an ingress queue passes XOFF, the switch sends a pause for that priority upstream; below XON it sends a zero-duration pause to resume.7 Two costs follow:
- Headroom. Packets already in flight when XOFF is sent must be absorbed, so each lossless priority on each port reserves buffer that scales with MTU, reaction time and above all cable length. With 9 or 12 MB shallow-buffer switches and links up to 300 m, Microsoft could afford two lossless classes even though PFC defines eight priorities.7
- Collateral damage. Pauses are per priority, not per flow, and propagate back toward sources: , victim flows, unfairness, congestion spreading.8 In a stress test on one of its test clusters Microsoft also hit a PFC deadlock caused by the interaction of flooding and pause propagation, and in production it saw NIC “pause storms” from a malfunctioning NIC.7
Losslessness is also never complete. In a Microsoft lab test, a switch set to drop 1 in 256 packets showed that the NIC’s go-back-0 recovery (restart the whole message on any loss) turned that into livelock: links at line rate, zero goodput, because a 4 MB message of 4,000 packets never finished. Moving to go-back-N fixed it.7
Do you need a lossless fabric at all?
Mittal et al. argued that PFC is an artifact of RoCE NIC design: with selective retransmission and a bounded in-flight window on the NIC (their IRN design, about 3–10% more NIC resources), RoCE without PFC outperformed RoCE with PFC by 6–83% in their scenarios.9 The transport (UET) takes the same position. It is explicitly designed for a best-effort network, and the specification says PFC should not be used anywhere in such a network because it violates the latency assumptions of UET’s congestion control.12 UET still defines optional link-level tools: credit-based flow control (CBFC) for classes that must be lossless, and link-layer retry (LLR), which replays frames lost to bit errors at the link rather than end to end.12 Its comparison of CBFC with PFC is a useful summary of the trade:12
| Credit-based (CBFC, InfiniBand-style) | Pause-based (PFC) |
|---|---|
| More lossless classes for the same buffer | Simpler XOFF/XON logic |
| Sender knows per-class credit, usable for scheduling and adaptive routing | Better sharing of burst buffer across ports |
| Underestimated cable delay only lowers throughput; PFC would overflow and drop | No messaging overhead when nothing is congested |
TCP/IP: the operating system on each CPU builds packets, copies buffers and handles acknowledgments. At 40 Gb/s that took 6% of a 32-core CPU to send and 12% to receive.
Flows A → R1 and B → R2 share the S1 → S2 link, in the same traffic class.
Traffic jams
Even a good network jams if too much traffic heads to one place. The classic case is : many chips send to one chip at the same moment. The last switch before it can’t pass everything on fast enough. AI training causes this a lot. Every chip finishes a step at about the same time, and they all send their results together.1514
A crowded switch can write a note on passing data: “I’m getting crowded.” The receiver tells the sender, and the sender slows down for a moment.8
Picking a lane
In a big network, there are many equally short routes between two switches. A switch picks one for each conversation, a bit like rolling a die. With thousands of small conversations, every route gets about the same share. But AI training sends a few huge transfers. Two of them can roll the same number and share one route, while another route sits empty. In one test, this bad luck wasted about 60% of the network.5
That second idea is called . The newest Ethernet rules for AI are built around it.12
Where congestion comes from
Congestion happens wherever the traffic arriving at a switch port exceeds the port’s speed. In training clusters the usual causes are (many senders to one receiver, as when several GPUs send to one at once), oversubscribed uplinks, and uneven use of the available paths. Switch buffers are small because switch bandwidth has grown faster than chip area and power allow buffer to grow,12 9 or 12 MB shared by all ports in Microsoft’s RDMA network,7 so a burst at line rate fills them quickly.
Congestion control: ECN and DCQCN
Flow control (credits or PFC) stops packets from being dropped, but it does it by stalling links. End-to-end congestion control tries to stop senders from overloading the network in the first place. On RoCE networks the best-known scheme is , built by Microsoft and Mellanox on the congestion-notification packets defined in the RoCEv2 standard:8
- A switch whose queue is longer than a threshold marks passing packets with (explicit congestion notification), a bit in the IP header.
- The receiving NIC sees the marks and sends a small congestion notification packet (CNP) back to the sender, at most one every 50 µs per flow.
- The sending NIC cuts its rate, then raises it again in steps if no more CNPs arrive.
In practice DCQCN is hard to tune for AI collectives. Meta deployed it for its 200 Gb/s clusters, found that settings which avoided PFC pauses hurt throughput, and ran its 400 Gb/s clusters with PFC alone, managing congestion instead by having the collective library admit traffic only when the receiver is ready.14
Load balancing: ECMP and its limits
A leaf–spine network built from 64-port switches has 32 equally short paths between any two leaves. The standard way to use them is (equal-cost multipath): each switch hashes header fields of a packet (addresses and ports) and picks a path from the result. Every packet of a flow gets the same hash, so a flow stays on one path and its packets stay in order.7 Al-Fares et al. pointed out the weakness in 2008: the hash ignores how big each flow is.3
That matters when traffic is a few . Hedera’s simulation of a 27,648-host fat tree found that hash collisions cut bisection bandwidth by 60.8% on average when each host sent one large flow at a time, but by only 2.5% with 1,000 simultaneous flows per host.5 AI training is the bad case: Meta reported ECMP performing poorly “due to the low flow entropy” of training traffic,14 and Alibaba found that the same flow hashed at three tiers in a row can line up badly at every one (hash polarization).15
Fixes
- More flows. Split each transfer over several connections so the hash has more to work with. Meta’s collective library does this (“QP scaling”).14
- Packet spraying. Send each packet of a flow down a different path. balances load almost perfectly, but packets arrive out of order, so the receiver must tolerate that.6
- Adaptive routing. Let switches pick each packet’s output port by current load. is built into current InfiniBand switches such as NVIDIA’s Quantum-2.20
- Topology. Fewer tiers mean fewer hashes in a row. Alibaba split its network into two independent planes and kept most jobs within one tier of rail switches partly for this reason.15
What Ultra Ethernet changes
The Ultra Ethernet transport is built for spraying. The sender’s NIC chooses an “entropy value” for each packet, which switches hash just as they hash ordinary ports, so the sender, not the switch, decides how to spread a flow over paths and can steer away from congested ones.12 Its delivery modes include reliable unordered delivery, so out-of-order packets are normal, and it defines two congestion-control algorithms: one run mainly by the sender using round-trip time and ECN, and one run by the receiver, which hands out permission to send.12 Switches may also “trim” a packet that doesn’t fit in the buffer, forwarding just its header so the receiver learns of the loss immediately.12
Congestion signals and DCQCN
Flow control (credits, PFC) bounds buffer occupancy hop by hop; congestion control bounds it end to end by adjusting injection. , the best-known scheme for RoCEv2, has three parts.8
- Congestion point (switch): marking on egress queue depth via RED, as in DCTCP (, , ).
- Notification point (receiver NIC): on a marked packet, send a CNP immediately if none was sent in the last µs, then at most one per µs ( in Microsoft’s deployment).
- Reaction point (sender NIC): on a CNP, , , . Without CNPs for , . Recovery is QCN-style: five rounds of fast recovery , then additive and hyper increase. Flows start at line rate; there is no slow start.
That is a lot of knobs, and they interact with PFC thresholds and with each other. At 200G, Meta found that tight ECN thresholds avoided PFC but cut collective throughput in corner cases, and that NIC settings which shaved about 3% off all-to-all completion time made PFC activity 2–3× worse. At 400G, performance degraded with default settings, partly from firmware changes. It now runs 400G RoCE with PFC only and controls congestion in the collective library: a receiver-driven admission scheme in which a sender transmits a chunk only after the receiver signals buffer space, with those clear-to-send packets given high priority in the switches.14 UET’s NSCC uses RTT and ECN together, and its receiver-credit algorithm (RCCC) explicitly targets incast; the specification expects RCCC to work best on non-blocking fat trees and to be complemented by NSCC where the fabric is oversubscribed.12
Load balancing with few, large flows
spreads flows by hashing; Al-Fares et al. noted that it does not account for flow bandwidth and that implementations were then limited to 8–16 ways.3 With flows hashed uniformly onto paths, the expected number of paths used is ; for that is 5.25, so about 34% of capacity sits idle while other links carry two or three flows (worked in “Under the hood”). Hedera quantifies the loss at 60.8% for a fat tree with one flow per host at a time, falling to 2.5% at 1,000 concurrent flows per host.5 Collective traffic sits at the bad end of that curve, and with synchronous collectives the job’s step time follows the most loaded link, not the average.
Production responses, roughly in order of how much they change:
- Overprovision. Meta doubled leaf uplink bandwidth (1:2 undersubscribed) as a stopgap against flow collisions, calling it expensive.14
- Pin paths. Meta’s path pinning by destination worked only with whole-rack job placement and no failures; fragmented placement degraded training by up to more than 30%.14
- Add entropy. QP scaling posts each message over several queue pairs, and enhanced ECMP hashes on the QP number as well; Meta used a scaling factor of 16 for LLM workloads.14
- Reduce hashing stages. Alibaba’s dual-plane design keeps each NIC port’s traffic in one plane, avoiding polarization at the aggregation layer, and its large tier-1 segments keep most traffic off the aggregation layer altogether.15
- Spray or adapt. removes collisions outright. Dixit et al. showed that in symmetric multi-rooted trees the parallel paths have similar queueing, so even TCP tolerates the reordering; asymmetry (for example a failed link) hurts it.6 picks ports by load in the switch and is built into current InfiniBand switches such as NVIDIA’s Quantum-2.20
Ultra Ethernet’s design
UET makes the endpoint responsible for path choice. Switches keep plain ECMP and are expected not to rewrite the entropy value; the sender’s congestion-management sublayer sets a per-packet entropy (carried in the UDP source port, or in a UET entropy header) and so chooses among paths, spreading a connection across all of them and moving off paths that report congestion.12 UET does not mandate a load-balancing algorithm.12 Supporting pieces:
- Delivery modes. Reliable unordered (RUD), reliable ordered (ROD), reliable unordered for idempotent operations, and unreliable unordered; an MPI library might use ordered delivery for headers and unordered for bulk payload.12
- Signals. ECN marked at dequeue rather than enqueue, so the mark reflects the queue the packet actually waited in.12
- Trimming. A switch that can’t buffer a packet may cut it to its header and forward that in a separate class. For a spraying transport, loss is hard to infer from out-of-order arrival, so the trimmed header gives the receiver a fast, explicit loss signal.12
ECN off: 8 senders at line rate into one port. The queue climbs to the pause threshold and stays there: the feeding links are paused 99% of the time, stalling any other traffic on them.
ECMP collision: both flows hash to spine 3, so each gets half the link; flow A finishes at t = 17.4. Packets stay in order. Three paths idle.
Build a network. Pick how many sockets each switch has, and two or three layers. Watch how many chips you can connect and how many switches it takes. Then slide “fewer cables going up” to the right. You save switches, but the green bar for “everyone sends to everyone” drops. Turn on rails to see each server’s cables fan out to eight different switches.
The builder computes a folded Clos (leaf–spine or three-tier fat tree) from the switch radix, tier count and leaf oversubscription, with 400 Gb/s per NIC and eight NICs per server. It reports endpoints, switches, cables, bisection bandwidth and worst-case hops, and estimates what share of line rate two traffic patterns get: an all-to-all among all GPUs, and a ring all-reduce. Try “Cloud” (3:1) and compare the two bars; then set radix 16 at 3:1 and toggle rails to see the ring recover.
The topology numbers are exact for the construction described in “Under the hood” (). The two bars are a leaf-uplink bound, , where is the share of a GPU’s traffic that must leave its leaf: uniform all-to-all (cross-rail hops ride the scale-up link, PXN-style), and one ring per GPU rank in placement order. They ignore congestion control, hashing and latency. The bottom panel hashes elephant flows onto equal-cost uplinks with a fixed hash; rehash to sample, and compare the slowest flow with packet spraying.
- Computers linked by cheap switches in three layers (2008)
- 27,648
- Fast home internet connections one switch chip could carry (2022)
- Half a million
- AI chips only one switch apart, using rails
- 1,024
- Network wasted by unlucky route picks (in a test)
- about 60%
- In 2008, researchers showed that cheap 48-socket switches in three layers could link 27,648 computers. Every one could talk at full speed.3
- One switch chip from 2022 moves as much data as half a million fast home internet connections at once.19
- Alibaba’s AI network uses rails so that 1,024 chips are only one switch apart. Most of its training jobs fit inside that group.15
- In a test where each computer sent one big transfer at a time, unlucky route picks wasted about 60% of the network.5
- Meta trained one of its Llama AI models on 16,000 chips, all joined by Ethernet.14
- Fat tree of 48-port switches
- 27,648 hosts
- RDMA vs TCP, 99th-percentile latency
- 90 vs 700 µs
- ECMP loss, 1 flow per host (k = 48)
- 60.8%
- Rail-only cost saving vs full-bisection Clos
- 38–77%
How big a network can you build?
For a non-blocking folded Clos, endpoints grow as for two tiers and for three, where is the switch radix (see “Switches in layers”). The table is pure arithmetic:
| Radix | 2 tiers: endpoints | 2 tiers: switches | 3 tiers: endpoints | 3 tiers: switches |
|---|---|---|---|---|
| 32 | 512 | 48 | 8,192 | 1,280 |
| 48 | 1,152 | 72 | 27,648 | 2,880 |
| 64 | 2,048 | 96 | 65,536 | 5,120 |
| 128 | 8,192 | 192 | 524,288 | 20,480 |
The three-tier figure matches Al-Fares et al.3 Real fabrics bend these numbers with oversubscription, rails, dual connections for redundancy and spare ports.
- Per-switch capacity: Jupiter chip (2012) → 51.2T chip
- 0.64 → 51.2 Tb/s
- HPN tier-1 ratio (128 × 200G : 60 × 400G)
- 1.067:1
- DCQCN CNP interval / α timer
- 50 / 55 µs
- ECMP loss: 1 vs 1,000 flows per host
- 60.8% vs 2.5%
Sizing, with real ratios
The pure-arithmetic sizes (, ) move once real constraints apply. Alibaba’s tier-1 switch is a 51.2 Tb/s chip run with 128 active plus 8 spare 200G ports down and 60 × 400G up: 25.6 Tb/s against 24 Tb/s, 1.067:1. Each host has eight 2 × 200G NICs, one per GPU, with the two ports going to different switches (dual-ToR), so a host connects to 16 tier-1 switches; 16 switches collectively serve 1,024 GPUs in a segment. Alibaba reports that 96.3% of its LLM training jobs use fewer than 1K GPUs and so fit in one segment.15 A dual-plane tier 2 then reaches about 15K GPUs per pod in two tiers, with 15:1 oversubscription to the core.15
Meta’s numbers frame the job side: it has described RoCE clusters of up to 32,000 GPUs, trained a large Llama 3 variant on 16,000 GPUs of a 24,000-GPU cluster, and notes that because of multi-dimensional parallelism the largest single collective typically spans only hundreds of GPUs even in jobs of tens of thousands.14 That is the quantitative case for oversubscribing high in the tree.
Lossless Ethernet in practice
| Quantity (Microsoft RoCEv2, 2016) | Value |
|---|---|
| Servers per ToR / link speed | 20–40 / 40 Gb/s |
| Cable runs: server–ToR, ToR–leaf, leaf–spine | ~2 m, 10–20 m, 200–300 m |
| ToR/leaf shared buffer | 9 or 12 MB |
| Lossless classes affordable (of 8 priorities) | 2 |
| 99th / 99.9th percentile latency, RDMA | 90 µs / ~200 µs (TCP 99th: 700 µs) |
All figures from Guo et al.7
- Bigger switches or more layers. Switches with more sockets mean fewer layers, so messages pass through fewer switches. But adding a layer lets the network grow far larger.
- Saving money or room for everyone. Fewer cables going up saves a lot of switches. That works when most traffic stays nearby. When every chip talks to every other, it becomes a jam.3
- Rails or simplicity. Rails put matching chips one switch apart. But chip 3 talking to chip 5 has to take a longer way around.15
- Never losing data or never stopping. Pausing a busy link means nothing gets lost. But it can hold up traffic that was nowhere near the jam. Letting some data get lost avoids that, but then the network cards must resend it quickly.9
- InfiniBand or Ethernet. InfiniBand was built for this job. Ethernet is what the rest of the datacenter already uses, and many companies sell it. It is now being reworked for AI.12
Tiers and radix
Each tier multiplies reach by roughly but adds two switch hops to the worst path, more switches per GPU ( for two tiers, for three) and more long links, which usually means more optical transceivers (see Optics). A higher-radix switch is the cleanest way to avoid a tier, which is why designers split 51.2 Tb/s chips into many slower ports rather than 64 × 800G: Alibaba’s HPN runs its tier-1 chip as 128 (plus 8 spare) × 200G down and 60 × 400G up.15
Oversubscription and placement
Oversubscribing saves the most hardware where the most links are, but it caps exactly the traffic that crosses the oversubscribed layer. Rings and pipeline hand-offs mostly stay local; all-to-all exchanges, used by mixture-of-experts models, cross the top of the network. The usual answer is non-blocking blocks that are big enough for most jobs, oversubscription above them, and a scheduler that keeps jobs inside a block when it can.1415
Rails
Rails cut hops and multiply the one-hop domain by eight for same-rank traffic, at the cost of making cross-rank traffic either climb or detour through the server’s scale-up link. If almost all traffic is same-rank, the spine can go entirely (rail-only), cutting network cost 38–77% compared with a full-bisection Clos, but jobs then depend on the scale-up domain to forward everything else.16
Lossless or not
PFC keeps RoCE simple at the NIC but brings head-of-line blocking, congestion spreading, and rare but serious failures such as deadlocks and pause storms.7 Smarter NIC loss recovery removes the need for PFC at a small hardware cost.9 InfiniBand’s credits need less buffer than PFC and are less sensitive to cable length, though credit-based flow control still shares its head-of-line blocking and congestion spreading, and it ties you to a fabric built for the purpose.129
Load balancing
Per-flow hashing keeps packets in order but collides; spraying balances but reorders; adaptive routing balances using live load information but, per packet, also reorders. Each pushes work to a different place: the network (overprovisioning), the switch (adaptive routing) or the NIC (reordering, path selection).56
How it goes wrong
- Failures unbalance the fabric. When a link or spine fails, the flows it carried get rehashed onto the survivors and collide; Meta saw jobs slow down from exactly this.14
- Fragmented placement. A job that only partly fills racks breaks assumptions that a static routing scheme relies on; Meta measured slowdowns of up to more than 30%.14
- One switch, many GPUs. A failed leaf takes every GPU behind it offline, and a synchronous job stops. Alibaba connects each NIC to two leaves for this reason.15
Tiers, radix and the optics bill
Per endpoint, a non-blocking two-tier fabric has switches and fabric links for endpoints, one per endpoint; a three-tier fabric has switches and two fabric links per endpoint. If every switch-to-switch link is optical, the transceiver count roughly doubles going from two tiers to three, on top of the extra hops. Meta’s leaf-to-spine links, for example, use single-mode fiber with 400G pluggable transceivers, while GPUs reach the leaf over copper within the rack.14 Radix is the lever that removes a tier, and lane speed and radix trade against each other on a fixed SerDes budget.19
Oversubscription is a scheduling contract
A u:d leaf caps cross-leaf throughput at u/d of line rate only for traffic that actually crosses; the delivered performance depends on what fraction of each GPU’s bytes must leave its block. Production designs make that fraction small by construction: Meta’s scheduler computes a minimum cut when splitting a job across AI Zones and assigns ranks accordingly;14 Alibaba places pipeline-parallel stages across pods so only that low-volume traffic meets the 15:1 core.15 The failure mode is a workload change, such as a mixture-of-experts model whose all-to-all crosses the top, landing on a fabric sized for rings.16
Rails versus generality
Rail-optimized fabrics keep a full Clos above the rails, so cross-rail traffic can still go up. Rail-only removes that layer: cheaper (38–77% below a full-bisection Clos in Wang et al.’s cost model), but cross-rail traffic must be forwarded through the scale-up domain, adding a “bandwidth tax” of extra bytes on those links and coupling network reachability to the scale-up fabric’s health.16
Lossless versus lossy, revisited
Lossless fabrics move the cost of congestion from retransmission into stalls: PFC’s headroom limits the number of lossless classes, and its pauses spread and occasionally deadlock.7 Lossy fabrics need endpoints with selective retransmission, bounded in-flight data and fast loss detection, which IRN showed costs a few percent of NIC resources9 and which UET builds in, with trimming to make loss explicit.12 Credit-based flow control sits in between: lossless without PFC’s sensitivity to cable delay, but more complex and, per UET, worse at sharing buffer across ports.12
InfiniBand versus Ethernet
Technically the gap has narrowed from both sides. InfiniBand brings credits, in-order lossless links, central routing, adaptive routing and in-network reduction as one integrated design.1020 Ethernet brings IP routing, multi-vendor switch silicon (in the two chips compared above, twice the radix at 400G) and shared operations with the rest of the datacenter, and is adding spraying, trimming, credits and link-level retry through Ultra Ethernet.12 The choice in practice also turns on supply, staff experience and how much tuning an operator is prepared to own; Meta’s account of moving from DCQCN to receiver-driven admission shows how much of that tuning lands on the operator with RoCE.14
Each see-saw leans toward the side of the bargain this design takes; level means it leaves the choice open. Tap one for details.
This part goes deeper, into the math, models and algorithms behind the chapter. It’s written for the Expert level.
1. Counting a folded Clos with oversubscription
Let every switch have radix , and let leaves have uplinks and downlinks (oversubscription ). The simulator sets .
Two tiers. Each spine has one port per leaf, so leaves and spines. Endpoints ; switches ; leaf–spine links ; equal-cost leaf-to-leaf paths ; worst-case 3 switch hops. For these reduce to , and .
Three tiers. Generalizing Al-Fares et al.’s construction: pods, each with leaves and aggregation switches (each aggregation switch has ports down, one per leaf in its pod, and up); , each with one port per pod. Endpoints ; switches ; pod-to-pod equal-cost paths ; worst case 5 switch hops. For : endpoints, switches, paths, matching the paper’s 27,648 hosts and 576 paths at .3
Rails. With eight NICs per server and no rails, a leaf holds whole servers and ports go unused (in the simulator). With rails, groups of eight leaves hold servers, one NIC per server per leaf, so all ports are used and same-rank GPUs on servers are one hop apart.
2. Why , and why fabrics settle for
Take a three-stage Clos with inputs per ingress switch. A new call from a free input on ingress switch A to a free output on egress switch C fails only if every middle switch is already busy toward A or toward C. At most middle links from A are busy (A’s other inputs) and at most toward C, so with middle switches at least one is free on both sides: strictly non-blocking. Clos’s 36-line example uses and .1 With the network is only rearrangeably non-blocking: any permutation can be routed at full rate, but adding a connection may require moving others. Packet networks don’t hold circuits, so they take the cheaper and rely on routing (hashing, spraying, adaptive choice) to come close to the rearranged optimum.3
m = 3: A’s two busy links and C’s two busy links cover all 3 middle switches, so A → C is blocked. Rearranging one existing call would free a path (m ≥ n).
3. Bisection and the leaf-uplink bound
In a folded Clos with leaves, any traffic between leaves must use leaf uplinks, and each leaf has of them for of injection. For a balanced cut between groups of leaves, the crossing capacity is , which is the simulator’s bisection readout. More generally, if a fraction of each GPU’s traffic must leave its leaf, the sustainable injection rate per GPU is bounded by
That is the model behind the simulator’s two bars. It assumes perfect load balancing over the uplinks and enough non-blocking capacity above the leaves (true for this construction, where aggregation switches are 1:1).
4. Ring all-reduce and all-to-all, in bytes and in placement
For an all-reduce of bytes over participants, at least one participant must send at least bytes, and a ring (reduce-scatter then all-gather around a ring) achieves that bound with each node talking only to its neighbors, contention-free on tree topologies.18 That neighbor-only property is what makes rings cheap on a fat tree:
- Ring, no rails. A leaf holds servers. With one ring per GPU rank, visiting servers in placement order, each of the 8 rings has exactly one edge leaving the leaf: .
- Ring, rails. Leaf holds rank of servers; its single ring leaves once: .
- All-to-all. Of the destinations outside a GPU’s own server, those under the same leaf are local. Without rails that is ; with rails, assuming cross-rail traffic first hops over the scale-up link to the destination’s rank (PXN), every GPU in the servers of the rail group is reachable at one hop, so .17 , close to 1 in any large network.
Worked example, , two tiers, 3:1 (, , ): . Rings: without rails, so throughput is . All-to-all: without rails, so about 34% of line rate; with rails and PXN , about 38%. Oversubscription is nearly free for rings and nearly fully paid by all-to-all. At and 3:1 (, ), a leaf holds only one whole server without rails, so and even the ring drops to a third; rails restore it to 100% by putting 12 servers per leaf group.
5. ECMP as balls into bins
Hash equal elephant flows independently and uniformly onto equal paths of capacity 1. A path carries nothing with probability , so
and aggregate throughput, if each loaded link runs full, equals the number of paths used. For : 5.25 of 8, so 34% of capacity is idle. The maximum load grows roughly like for , so some link almost always carries two or three flows; those flows get 1/2 or 1/3 of line rate, and a synchronous collective finishes at their pace. With the relative imbalance shrinks, which is Hedera’s 60.8% versus 2.5% result: many small flows per host average out, one big one does not.5 Hashing at several tiers in a row with the same inputs makes it worse, because the choices are correlated rather than independent (polarization).15 Packet spraying turns the problem into tiny balls, and balance becomes nearly perfect, at the price of reordering.6
6. DCQCN as a control loop
DCQCN is multiplicative decrease with a DCTCP-style estimate of the marking fraction: tracks the recent probability that a CNP window saw a mark, and each CNP cuts the rate by .8 Its feedback is rate-limited to one CNP per flow per 50 µs, and its increase is timer- and byte-counter-driven rather than ACK-clocked, which is why its behavior under synchronized incast depends so much on parameters: in Meta’s experiments, ECN settings tight enough to avoid PFC cut collective throughput in some cases, and relaxed settings that trimmed completion time about 3% made PFC 2–3× worse.14 UET’s NSCC instead combines RTT and ECN at the sender, and RCCC hands out receiver credits, targeting the many-to-one case directly.12
Q1A two-tier leaf–spine network uses 64-port switches with half of each leaf’s ports going down and half going up. How many servers’ NIC ports can it connect at most?
Q2A leaf switch has 48 ports to servers and 16 uplinks, all at the same speed. If every server sends to servers under other leaves at once, what share of its NIC speed can each get?
Q3Why do RoCE deployments usually turn on priority flow control (PFC)?
Q4Oversubscription hurts an all-to-all exchange much more than a ring all-reduce. Why?
Sources
Show Hide 20 sources
- A Study of Non-Blocking Switching NetworksThree-stage switching arrays that need fewer crosspoints than one big square switch; the 36-input example with six 6 × 11 input switches, eleven 6 × 6 middle switches and 1,188 crosspoints; why the middle stage needs 2n − 1 switches to never block.
- Fat-Trees: Universal Networks for Hardware-Efficient SupercomputingProcessors at the leaves, switches at the internal nodes, and more wires toward the root so bandwidth grows going up; for a given hardware volume no network is much better than a fat-tree.
- A Scalable, Commodity Data Center Network ArchitectureDefinition of oversubscription and typical 2.5:1 to 8:1 designs; the k-ary fat-tree (k pods, k/2 + k/2 switches per pod, (k/2)² cores, k³/4 hosts); 48-port switches give 27,648 hosts and 576 equal-cost paths; rearrangeably non-blocking; ECMP’s static per-flow splitting ignores flow size.
- Jupiter Rising: A Decade of Clos Topologies and Centralized Control in Google’s Datacenter NetworkFive generations of multi-stage Clos fabrics built from merchant switch silicon; Jupiter (2012) built from 16 × 40G chips reaches 1.3 Pb/s of bisection bandwidth; capacity grew 100× in ten years.
- Hedera: Dynamic Flow Scheduling for Data Center NetworksLarge long-lived flows collide on ECMP hashes; in a k = 48 fat-tree, one flow per host at a time loses 60.8% of bisection bandwidth on average, while 1,000 parallel flows per host lose only 2.5%.
- On the Impact of Packet Spraying in Data Center NetworksPer-flow ECMP can leave load imbalanced; spraying packets of one flow over all equal-cost paths balances load, and the symmetry of multi-rooted trees keeps reordering tolerable, until a failure breaks the symmetry.
- RDMA over Commodity Ethernet at ScaleRoCEv2 carries the RDMA transport in UDP (destination port 4791, per-QP source port for ECMP); PFC pause, XOFF/XON and headroom; only two lossless classes fit in 9 or 12 MB switch buffers; a PFC deadlock in a test cluster and NIC pause storms in production; go-back-0 livelock in a lab test; RDMA 99th-percentile latency 90 µs vs 700 µs for TCP.
- Congestion Control for Large-Scale RDMA DeploymentsRDMA puts the transport on the NIC and bypasses the host stack; InfiniBand’s link layer uses hop-by-hop credit-based flow control; PFC causes head-of-line blocking and unfairness; DCQCN: switches ECN-mark, receivers send CNPs at most every 50 µs, senders cut rate by α/2.
- Revisiting Network Support for RDMA (extended version)PFC brings head-of-line blocking, congestion spreading and deadlocks; the need for it is an artifact of RoCE NIC loss recovery; an improved NIC (IRN) without PFC beats RoCE with PFC by 6–83% at 3–10% extra NIC resources.
- RDMA-aware Networks Programming Guide§3.1: the InfiniBand link layer has credit-based flow control and virtual lanes and guarantees strong ordering within a VL along a path; §3.5: the subnet manager develops a routing table, switches implement link-layer flow control to prevent packet dropping and support adaptive routing.
- NVIDIA SMOne subnet manager runs per InfiniBand subnet; OpenSM scans and initializes the fabric, sweeps for changes, and recalculates routes when a port goes down or a link is added.
- Ultra Ethernet Specification v1.0UET addresses RoCEv2’s shortcomings; scales to millions of endpoints; packet spraying steered by entropy values over ECMP; NSCC (RTT + ECN, sender) and RCCC (receiver) congestion control; packet trimming; PFC should not be used in best-effort networks; credit-based flow control and link-layer retry; delivery modes including reliable unordered (RUD).
- Ultra Ethernet Specification (download page)Public releases: 1.0 (June 11, 2025), 1.0.1 (September 5, 2025), 1.0.2 (January 28, 2026), 1.0.3 (July 16, 2026).
- RDMA over Ethernet for Distributed AI Training at Meta ScaleSeparate back-end RoCE network; one NIC per GPU, eight per host; a non-blocking two-stage Clos “AI Zone” with an oversubscribed aggregation layer above it; a 24,000-GPU cluster; ECMP’s low entropy, path pinning, QP scaling; doubled uplinks as a stopgap; DCQCN dropped at 400G in favor of PFC plus receiver-driven admission in the collective library.
- Alibaba HPN: A Data Center Network for Large Language Model TrainingLLM training produces few, periodic, bursty 400 Gb/s flows per host, so ECMP suffers hash polarization; 8 × 400 Gb/s per host; rail-optimized dual-ToR tier 1 puts 1,024 GPUs one switch apart (8× more than without rails); 51.2 Tb/s switch with 128 × 200G down and 60 × 400G up (1.067:1); 15K GPUs per two-tier pod; 15:1 aggregation-to-core.
- Rail-only: A Low-Cost High-Performance Network for Training LLMs with Trillion ParametersA rail is the set of GPUs with the same local rank in different high-bandwidth domains; rail-optimized networks put each rail under the same switches; LLM traffic mostly stays within a rail, so removing the spine matches rail-optimized performance at 38–77% lower network cost than a full-bisection Clos; cross-rail traffic is forwarded through the high-bandwidth domain, costing 8.2–11.2% on MoE all-to-all.
- Doubling all2all Performance with NVIDIA Collective Communication Library 2.12Rail-optimized topology (NIC 0 of every server to leaf 0, NIC 1 to leaf 1, …); PXN moves data over NVLink to a GPU next to the right NIC; more than 2× faster all2all.
- Bandwidth Optimal All-reduce Algorithms for Clusters of WorkstationsSome process must send at least 2(N − 1)/N × the data in an all-reduce; a ring algorithm achieves this bound contention-free on tree topologies.
- Broadcom Ships Tomahawk 5, Industry’s Highest Bandwidth Switch Chip to Accelerate AI/ML WorkloadsAnnounced August 16, 2022: 51.2 Tb/s of Ethernet switching in one monolithic 5 nm die; 64 ports of 800GbE or 256 ports of 200GbE; 512 × 100G PAM4 SerDes.
- QM97XX 1U NDR 400Gb/s InfiniBand Switch Systems User Manual: Introduction64 ports of NDR 400 Gb/s InfiniBand in 1U; 51.2 Tb/s aggregated bidirectional throughput; up to 128 ports of 200 Gb/s with port split; adaptive routing and SHARP in-network reduction.