Systems · Chapter 4 of 7 · Between the boxes

Scale-out networking

Teaching a big AI takes thousands of chips in many racks, joined by cables, like a very fast internet inside one building.

Scale-out networks connect servers through layers of switches. RDMA lets one machine write straight into another’s memory. InfiniBand and Ethernet compete here, and the topology decides how much bandwidth can cross the whole network at once.

Fat-tree, rail-optimized and other topologies, bisection bandwidth and oversubscription, RDMA and RoCE, congestion control and load balancing, InfiniBand versus Ethernet, and the Ultra Ethernet Consortium’s specifications.

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. A typical AI server today has eight accelerators and eight NICs, one per accelerator, each at 400 gigabits per second (Gb/s). The traffic is unusual by data-center standards: a few very large, periodic, bursty flows per server rather than millions of small ones.

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. Each accelerator gets its own NIC, typically 400 Gb/s, so an eight-GPU server injects 3.2 Tb/s.

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. 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.

4 leaves + 2 spines, 16 linksspineleaf12345678chips (each with a NIC)10,000 chipsdirect: n(n − 1)/249,995,000 cables9,999 ports per chip64-port switches30,032 cables783 switches, 3 tiers≤ 5 switch hops
Wiring

3 tiers of 64-port switches: 783 switches, 30,032 cables, at most 5 switch hops. Cables grow with n, not n².

Direct wiring against a network of switches. The drawing is the smallest leaf–spine (4-port switches, 8 chips); the counts use 64-port switches and no oversubscription.Share freely with credit: ‘Figure from chipfieldguide.com’

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.

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. AI clusters usually don’t, because all their chips really do talk at once.

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.

Radix and reach

A switch’s is its number of ports. One switch with radix kk can connect at most kk 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 kk:

  • each spine has kk ports, one per leaf, so there can be kk leaves;
  • each leaf has k/2k/2 ports down, so the network connects k×k/2=k2/2k \times k/2 = k^2/2 NIC ports;
  • any path is leaf → spine → leaf: at most three switch hops;
  • between any two leaves there are k/2k/2 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. 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. In 2008, Al-Fares, Loukissas and Vahdat showed how to build one from identical, cheap Ethernet switches. Their three-tier kk-ary fat tree has kk “pods,” each with k/2k/2 edge (leaf) and k/2k/2 aggregation switches, plus (k/2)2(k/2)^2 core switches on top, and supports k3/4k^3/4 hosts. With 48-port switches that is 27,648 hosts. 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.

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 2,048×400/2≈4102{,}048 \times 400 / 2 \approx 410 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. 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.

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.” 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. 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. 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.

Folded Clos, counted

Clos’s 1953 construction has an ingress stage of rr switches with nn inputs each, mm 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 n−1n - 1 of the ingress switch’s middle links and n−1n - 1 of the egress switch’s are busy, on disjoint middle switches; so m=2n−1m = 2n - 1 middle switches guarantee a free path (his 36-port example: n=6n = 6, m=11m = 11). Data-center fabrics fold the network so the ingress and egress switches are the same leaves, and settle for m=nm = n (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.

For radix kk and a non-oversubscribed fabric:

2 tiers (leaf–spine)3 tiers (kk-ary fat tree)
Endpointsk2/2k^2/2k3/4k^3/4
Switchesk+k/2=3k/2k + k/2 = 3k/2k2+(k/2)2=5k2/4k^2 + (k/2)^2 = 5k^2/4
Switches per endpoint3/k3/k5/k5/k
Worst-case switch hops35
Equal-cost paths, far endpointsk/2k/2(k/2)2(k/2)^2

The three-tier row is Al-Fares et al.’s construction: kk pods of k/2k/2 edge and k/2k/2 aggregation switches, (k/2)2(k/2)^2 cores, k3/4k^3/4 hosts; for k=48k = 48, 27,648 hosts and 576 equal-cost paths between hosts in different pods. 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. 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.

Radix is the master variable

Because reach goes as k2k^2 or k3k^3, 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 (3/k3/k versus 5/k5/k) 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. 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.

Bisection bandwidth and where to oversubscribe

is the minimum, over all balanced cuts, of the capacity crossing the cut. For a folded Clos with NN endpoints at rate BB and leaf ratio d:ud{:}u (downlinks to uplinks, all same speed), cuts between leaves are bounded by the leaf uplinks: NB/2×min⁡(1,u/d)NB/2 \times \min(1, u/d). 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.

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. 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.

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. Two things change:

  • Locality. A rail group of eight leaves with dd downlinks each serves dd servers, and all same-rank traffic among them is one hop. Without rails the same eight leaves serve dd servers too, but only d/8d/8 per leaf, so the one-hop set is 8× smaller.
  • 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. NCCL calls the second option PXN; NVIDIA measured all-to-all running more than twice as fast with it.

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.

coreagg.edgehostspod 1pod 2pod 3pod 4AB5 switch hops · 4 equal-cost paths
Tiers

A → B (another pod): 5 switch hops, 4 equal-cost paths (showing 1).

The smallest folded Clos networks (4-port switches). Tap a host to route from A; “Next path” steps through the equal-cost paths.Share freely with credit: ‘Figure from chipfieldguide.com’
spinesleavesserver 112A4server 21234server 312B4server 412343 switch hops · same rank
Wiring

Without rails each server sits on its own leaf, so any server-to-server traffic goes leaf → spine → leaf: three hops.

Drawn with 4 GPUs per server (real servers have 8). Tap a GPU to route from A; switch between one-leaf-per-server and rail-optimized wiring.Share freely with credit: ‘Figure from chipfieldguide.com’

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.

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.”

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.

Both kinds are used to train the biggest models. A newer version, Ultra Ethernet, aims to work without the stop sign at all.

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%. 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. In Microsoft’s production measurements, RDMA’s 99th-percentile latency was 90 microseconds (µs) against 700 µs for TCP on the same network. For AI, NICs can also read and write GPU memory directly, so GPU-to-GPU traffic skips host memory entirely.

The catch is that the NIC’s transport is simple. It was designed assuming the network almost never drops packets. 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. Packets are not dropped for lack of buffer space, and the link delivers them in order, so the transport layer stays simple. 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.

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.

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. 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.

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.

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.

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. 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. 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. 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. Routing is computed centrally by a subnet manager that programs paths into every switch and recomputes them on link changes. Current InfiniBand switches add adaptive routing and in-network aggregation and reduction (SHARP), so part of a collective’s arithmetic happens in the switches.

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. 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. 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.
  • Collateral damage. Pauses are per priority, not per flow, and propagate back toward sources: , victim flows, unfairness, congestion spreading. 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.

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.

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. 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. 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. Its comparison of CBFC with PFC is a useful summary of the trade:

Credit-based (CBFC, InfiniBand-style)Pause-based (PFC)
More lossless classes for the same bufferSimpler XOFF/XON logic
Sender knows per-class credit, usable for scheduling and adaptive routingBetter sharing of burst buffer across ports
Underestimated cable delay only lowers throughput; PFC would overflow and dropNo messaging overhead when nothing is congested
senderreceivermemoryCPUOS stackNICCPUOS stackmemoryNIC6% CPU (send)12% CPU (receive)p99 latency700 µs
Transport

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.

The data path for TCP and for RDMA. CPU and latency numbers are Microsoft’s measurements on its data-center network.Share freely with credit: ‘Figure from chipfieldguide.com’
1. Normal trafficABR1R2S1S2queue → R1
Full-buffer policy
1 / 4

Flows A → R1 and B → R2 share the S1 → S2 link, in the same traffic class.

One congested receiver, one innocent flow, three ways to handle a full buffer. Step through each and watch flow B.Share freely with credit: ‘Figure from chipfieldguide.com’

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.

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.

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.

That second idea is called . The newest Ethernet rules for AI are built around it.

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, 9 or 12 MB shared by all ports in Microsoft’s RDMA network, 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:

  1. A switch whose queue is longer than a threshold marks passing packets with (explicit congestion notification), a bit in the IP header.
  2. 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.
  3. 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.

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. Al-Fares et al. pointed out the weakness in 2008: the hash ignores how big each flow is.

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. AI training is the bad case: Meta reported ECMP performing poorly “due to the low flow entropy” of training traffic, and Alibaba found that the same flow hashed at three tiers in a row can line up badly at every one (hash polarization).

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”).
  • 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.
  • 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.
  • 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.

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. 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. 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.

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.

  • Congestion point (switch): marking on egress queue depth via RED, as in DCTCP (KminK_{\mathrm{min}}, KmaxK_{\mathrm{max}}, PmaxP_{\mathrm{max}}).
  • Notification point (receiver NIC): on a marked packet, send a CNP immediately if none was sent in the last NN µs, then at most one per NN µs (N=50 μsN = 50\,\mu\mathrm{s} in Microsoft’s deployment).
  • Reaction point (sender NIC): on a CNP, RT←RCR_{\mathrm{T}} \leftarrow R_{\mathrm{C}}, RC←RC(1−α/2)R_{\mathrm{C}} \leftarrow R_{\mathrm{C}}(1 - \alpha/2), α←(1−g)α+g\alpha \leftarrow (1 - g)\alpha + g. Without CNPs for K=55 μsK = 55\,\mu\mathrm{s}, α←(1−g)α\alpha \leftarrow (1 - g)\alpha. Recovery is QCN-style: five rounds of fast recovery RC←(RT+RC)/2R_{\mathrm{C}} \leftarrow (R_{\mathrm{T}} + R_{\mathrm{C}})/2, 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. 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.

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. With FF flows hashed uniformly onto PP paths, the expected number of paths used is P(1−(1−1/P)F)P\left(1 - (1 - 1/P)^F\right); for F=P=8F = P = 8 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 k=48k = 48 fat tree with one flow per host at a time, falling to 2.5% at 1,000 concurrent flows per host. 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.
  • 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%.
  • 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.
  • 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.
  • 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. picks ports by load in the switch and is built into current InfiniBand switches such as NVIDIA’s Quantum-2.

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. UET does not mandate a load-balancing algorithm. 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.
  • Signals. ECN marked at dequeue rather than enqueue, so the mark reflects the queue the packet actually waited in.
  • 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.
Rswitchqueue at R’s port (µs)pause threshold (PFC)red ticks: links pausedR’s link use100%00200400600800 µs
ECN + DCQCN

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.

Incast into one switch port, with and without ECN-based congestion control (DCQCN). Illustrative toy model: flows start at line rate, one CNP per 50 µs per flow at most.Share freely with credit: ‘Figure from chipfieldguide.com’
1idle2idle3A1–8, B1–84idle4 equal-cost spinesL1L2ABA′B′A′ arrivals12345678A done at t =17.4
Load balancing

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.

Per-flow ECMP (here the hash happened to collide) against packet spraying, for two 8-packet flows over four spines. Delays are illustrative.Share freely with credit: ‘Figure from chipfieldguide.com’

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” (uplinks=round(k/(1+r))\text{uplinks} = \mathrm{round}(k/(1 + r))). The two bars are a leaf-uplink bound, min⁡(1,(u/d)/f)\min(1, (u/d)/f), where ff 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.

Loading simulation…
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.
  • One switch chip from 2022 moves as much data as half a million fast home internet connections at once.
  • 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.
  • In a test where each computer sent one big transfer at a time, unlucky route picks wasted about 60% of the network.
  • Meta trained one of its Llama AI models on 16,000 chips, all joined by Ethernet.
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 k2/2k^2/2 for two tiers and k3/4k^3/4 for three, where kk is the switch radix (see “Switches in layers”). The table is pure arithmetic:

Radix kk2 tiers: endpoints2 tiers: switches3 tiers: endpoints3 tiers: switches
32512488,1921,280
481,1527227,6482,880
642,0489665,5365,120
1288,192192524,28820,480

The k=48k = 48 three-tier figure matches Al-Fares et al. 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 (k2/2k^2/2, k3/4k^3/4) 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. A dual-plane tier 2 then reaches about 15K GPUs per pod in two tiers, with 15:1 oversubscription to the core.

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. That is the quantitative case for oversubscribing high in the tree.

Lossless Ethernet in practice

Quantity (Microsoft RoCEv2, 2016)Value
Servers per ToR / link speed20–40 / 40 Gb/s
Cable runs: server–ToR, ToR–leaf, leaf–spine~2 m, 10–20 m, 200–300 m
ToR/leaf shared buffer9 or 12 MB
Lossless classes affordable (of 8 priorities)2
99th / 99.9th percentile latency, RDMA90 µs / ~200 µs (TCP 99th: 700 µs)

All figures from Guo et al.

  • 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.
  • Rails or simplicity. Rails put matching chips one switch apart. But chip 3 talking to chip 5 has to take a longer way around.
  • 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.
  • 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.

Tiers and radix

Each tier multiplies reach by roughly k/2k/2 but adds two switch hops to the worst path, more switches per GPU (3/k3/k for two tiers, 5/k5/k 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.

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.

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.

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. Smarter NIC loss recovery removes the need for PFC at a small hardware cost. 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.

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).

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.
  • 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%.
  • 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.

Tiers, radix and the optics bill

Per endpoint, a non-blocking two-tier fabric has 3/k3/k switches and k/2×kk/2 \times k fabric links for k2/2k^2/2 endpoints, one per endpoint; a three-tier fabric has 5/k5/k 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. Radix is the lever that removes a tier, and lane speed and radix trade against each other on a fixed SerDes budget.

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; Alibaba places pipeline-parallel stages across pods so only that low-volume traffic meets the 15:1 core. 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.

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.

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. Lossy fabrics need endpoints with selective retransmission, bounded in-flight data and fast loss detection, which IRN showed costs a few percent of NIC resources and which UET builds in, with trimming to make loss explicit. 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.

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. 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. 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.

fewer tiersmore tiersoversubscribednon-blockingrail-optimizedserver per leaflosslesslossy + resendper-flow pathsper-packet spray
Design

Each see-saw leans toward the side of the bargain this design takes; level means it leaves the choice open. Tap one for details.

Five trade-offs, and which side three designs take. Tilts are qualitative, based on the chapter’s sources.Share freely with credit: ‘Figure from chipfieldguide.com’

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 kk, and let leaves have uu uplinks and d=k−ud = k - u downlinks (oversubscription r=d/ur = d/u). The simulator sets u=round(k/(1+r))u = \mathrm{round}(k/(1 + r)).

Two tiers. Each spine has one port per leaf, so L=kL = k leaves and S=uS = u spines. Endpoints N=kdN = kd; switches k+uk + u; leaf–spine links kuku; equal-cost leaf-to-leaf paths uu; worst-case 3 switch hops. For r=1r = 1 these reduce to k2/2k^2/2, 3k/23k/2 and k/2k/2.

Three tiers. Generalizing Al-Fares et al.’s construction: kk pods, each with k/2k/2 leaves and uu aggregation switches (each aggregation switch has k/2k/2 ports down, one per leaf in its pod, and k/2k/2 up); cores=uk/2\text{cores} = uk/2, each with one port per pod. Endpoints N=k⋅(k/2)⋅dN = k \cdot (k/2) \cdot d; switches k⋅k/2+ku+uk/2k \cdot k/2 + ku + uk/2; pod-to-pod equal-cost paths uk/2uk/2; worst case 5 switch hops. For r=1r = 1: k3/4k^3/4 endpoints, 5k2/45k^2/4 switches, (k/2)2(k/2)^2 paths, matching the paper’s 27,648 hosts and 576 paths at k=48k = 48.

Rails. With eight NICs per server and no rails, a leaf holds ⌊d/8⌋\lfloor d/8 \rfloor whole servers and d mod 8d \bmod 8 ports go unused (in the simulator). With rails, groups of eight leaves hold dd servers, one NIC per server per leaf, so all dd ports are used and same-rank GPUs on dd servers are one hop apart.

2. Why 2n−12n - 1, and why fabrics settle for nn

Take a three-stage Clos with nn 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 n−1n - 1 middle links from A are busy (A’s other inputs) and at most n−1n - 1 toward C, so with 2n−12n - 1 middle switches at least one is free on both sides: strictly non-blocking. Clos’s 36-line example uses n=6n = 6 and m=11m = 11. With m=nm = n 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 m=nm = n and rely on routing (hashing, spraying, adaptive choice) to come close to the rearranged optimum.

A → C blockedAXCYM1A→Y M2A→Y X→CM3X→Cingressegressn = 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).

A three-stage Clos network, n = 3, in its worst case for a new call A → C. Slide the number of middle switches m; when blocked, try rearranging.Share freely with credit: ‘Figure from chipfieldguide.com’

3. Bisection and the leaf-uplink bound

In a folded Clos with d:ud{:}u leaves, any traffic between leaves must use leaf uplinks, and each leaf has uBuB of them for dBdB of injection. For a balanced cut between groups of leaves, the crossing capacity is NB/2×min⁡(1,u/d)NB/2 \times \min(1, u/d), which is the simulator’s bisection readout. More generally, if a fraction ff of each GPU’s traffic must leave its leaf, the sustainable injection rate per GPU is bounded by

throughputB≤min⁡ ⁣(1,u/df)\frac{\text{throughput}}{B} \le \min\!\left(1, \frac{u/d}{f}\right)

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 XX bytes over NN participants, at least one participant must send at least 2(N−1)/N⋅X2(N - 1)/N \cdot X 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. That neighbor-only property is what makes rings cheap on a fat tree:

  • Ring, no rails. A leaf holds s=⌊d/8⌋s = \lfloor d/8 \rfloor servers. With one ring per GPU rank, visiting servers in placement order, each of the 8 rings has exactly one edge leaving the leaf: f=8/(8s)=1/sf = 8/(8s) = 1/s.
  • Ring, rails. Leaf ii holds rank ii of dd servers; its single ring leaves once: f=1/df = 1/d.
  • All-to-all. Of the N−8N - 8 destinations outside a GPU’s own server, those under the same leaf are local. Without rails that is 8s−88s - 8; with rails, assuming cross-rail traffic first hops over the scale-up link to the destination’s rank (PXN), every GPU in the dd servers of the rail group is reachable at one hop, so 8d−88d - 8. f=(N−local−8)/(N−8)f = (N - \text{local} - 8)/(N - 8), close to 1 in any large network.

Worked example, k=64k = 64, two tiers, 3:1 (d=48d = 48, u=16u = 16, N=3,072N = 3{,}072): u/d=1/3u/d = 1/3. Rings: f=1/6f = 1/6 without rails, so throughput is min⁡(1,(1/3)/(1/6))=100%\min(1, (1/3)/(1/6)) = 100\%. All-to-all: f≈0.99f \approx 0.99 without rails, so about 34% of line rate; with rails and PXN f≈0.88f \approx 0.88, about 38%. Oversubscription is nearly free for rings and nearly fully paid by all-to-all. At k=16k = 16 and 3:1 (d=12d = 12, u=4u = 4), a leaf holds only one whole server without rails, so f=1f = 1 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 FF equal elephant flows independently and uniformly onto PP equal paths of capacity 1. A path carries nothing with probability (1−1/P)F(1 - 1/P)^F, so

E[paths used]=P(1−(1−1P)F)E[\text{paths used}] = P\left(1 - \left(1 - \frac{1}{P}\right)^{F}\right)

and aggregate throughput, if each loaded link runs full, equals the number of paths used. For F=P=8F = P = 8: 5.25 of 8, so 34% of capacity is idle. The maximum load grows roughly like ln⁡P/ln⁡ln⁡P\ln P / \ln \ln P for F=PF = P, 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 F≫PF \gg P 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. Hashing at several tiers in a row with the same inputs makes it worse, because the choices are correlated rather than independent (polarization). Packet spraying turns the problem into F⋅(packets)F \cdot (\text{packets}) tiny balls, and balance becomes nearly perfect, at the price of reordering.

6. DCQCN as a control loop

DCQCN is multiplicative decrease with a DCTCP-style estimate of the marking fraction: α\alpha tracks the recent probability that a CNP window saw a mark, and each CNP cuts the rate by α/2\alpha/2. 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. UET’s NSCC instead combines RTT and ECN at the sender, and RCCC hands out receiver credits, targeting the many-to-one case directly.

Novice · 0 of 4 correct
  1. 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?

  2. 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?

  3. Q3Why do RoCE deployments usually turn on priority flow control (PFC)?

  4. Q4Oversubscription hurts an all-to-all exchange much more than a ring all-reduce. Why?

Sources

Show Hide 20 sources
  1. A Study of Non-Blocking Switching NetworksCharles Clos · Bell System Technical Journal 32(2), via Internet Archive · 1953Three-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.
  2. Fat-Trees: Universal Networks for Hardware-Efficient SupercomputingCharles E. Leiserson · IEEE Transactions on Computers C-34(10), copy on the author’s MIT course page (6.896) · 1985Processors 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.
  3. A Scalable, Commodity Data Center Network ArchitectureMohammad Al-Fares, Alexander Loukissas and Amin Vahdat · ACM SIGCOMM 2008, author copy (Amin Vahdat, UC San Diego) · 2008Definition 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.
  4. Jupiter Rising: A Decade of Clos Topologies and Centralized Control in Google’s Datacenter NetworkArjun Singh et al. (Google) · ACM SIGCOMM 2015, Google Research · 2015Five 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.
  5. Hedera: Dynamic Flow Scheduling for Data Center NetworksMohammad Al-Fares, Sivasankar Radhakrishnan, Barath Raghavan, Nelson Huang and Amin Vahdat · USENIX NSDI 2010, author copy · 2010Large 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%.
  6. On the Impact of Packet Spraying in Data Center NetworksAdvait Dixit, Pawan Prakash, Y. Charlie Hu and Ramana Rao Kompella · IEEE INFOCOM 2013, author copy at Purdue University · 2013Per-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.
  7. RDMA over Commodity Ethernet at ScaleChuanxiong Guo, Haitao Wu, Zhong Deng, Gaurav Soni, Jianxi Ye, Jitendra Padhye and Marina Lipshteyn (Microsoft) · ACM SIGCOMM 2016, author copy at Microsoft Research · 2016RoCEv2 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.
  8. Congestion Control for Large-Scale RDMA DeploymentsYibo Zhu et al. (Microsoft, Mellanox) · ACM SIGCOMM 2015, author copy at Microsoft Research · 2015RDMA 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.
  9. Revisiting Network Support for RDMA (extended version)Radhika Mittal, Alexander Shpiner, Aurojit Panda, Eitan Zahavi, Arvind Krishnamurthy, Sylvia Ratnasamy and Scott Shenker · arXiv:1806.08159 (SIGCOMM 2018) · 2018PFC 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.
  10. RDMA-aware Networks Programming GuideNVIDIA · NVIDIA DOCA SDK documentation v3.5.0 · 2026§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.
  11. NVIDIA SMNVIDIA · NVIDIA DOCA SDK documentation v3.5.0 · 2026One 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.
  12. Ultra Ethernet Specification v1.0Ultra Ethernet Consortium · Ultra Ethernet Consortium · 2025UET 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).
  13. Ultra Ethernet Specification (download page)Ultra Ethernet Consortium · Ultra Ethernet ConsortiumPublic releases: 1.0 (June 11, 2025), 1.0.1 (September 5, 2025), 1.0.2 (January 28, 2026), 1.0.3 (July 16, 2026).
  14. RDMA over Ethernet for Distributed AI Training at Meta ScaleAdithya Gangidi et al. (Meta) · ACM SIGCOMM 2024, Meta Engineering · 2024Separate 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.
  15. Alibaba HPN: A Data Center Network for Large Language Model TrainingKun Qian et al. (Alibaba Cloud) · ACM SIGCOMM 2024, author copy · 2024LLM 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.
  16. Rail-only: A Low-Cost High-Performance Network for Training LLMs with Trillion ParametersWeiyang Wang, Manya Ghobadi, Kayvon Shakeri, Ying Zhang and Naader Hasani · arXiv:2307.12169 · 2023A 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.
  17. Doubling all2all Performance with NVIDIA Collective Communication Library 2.12Karthik Mandakolathur and Sylvain Jeaugey · NVIDIA Technical Blog · 2022Rail-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.
  18. Bandwidth Optimal All-reduce Algorithms for Clusters of WorkstationsPitch Patarasuk and Xin Yuan · Journal of Parallel and Distributed Computing, author copy at Florida State University · 2009Some 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.
  19. Broadcom Ships Tomahawk 5, Industry’s Highest Bandwidth Switch Chip to Accelerate AI/ML WorkloadsBroadcom Inc. · Broadcom (news release) · 2022Announced 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.
  20. QM97XX 1U NDR 400Gb/s InfiniBand Switch Systems User Manual: IntroductionNVIDIA · NVIDIA Networking documentation64 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.