A shared-memory multiprocessor stops growing at a few hundred cores: keeping one coherent address space across more is too hard and too expensive. The biggest computers in the world take the other road. They are thousands of ordinary computers, each with its own memory and its own operating system, that cooperate by sending each other messages over a fast network. A web search, a weather forecast and the training of a large AI model all run on machines like that.
This chapter is about those multicomputers: how programs on them communicate, what a message costs, how the network is wired, and what today's largest systems look like.
Message passing
In a multicomputer, no processor can load from another node's memory. If node 3 needs data held by node 7, node 7 has to send it and node 3 has to receive it. Communication is explicit, in the program, instead of hidden behind loads and stores.
The standard interface for scientific computing is MPI (Message Passing Interface), first published in 1994 and implemented for every supercomputer since. Each process of a job has a rank, and ranks exchange messages:
#include <mpi.h>
int main(int argc, char **argv) {
int rank;
double x = 0;
MPI_Init(&argc, &argv);
MPI_Comm_rank(MPI_COMM_WORLD, &rank);
if (rank == 0) {
x = 3.14;
MPI_Send(&x, 1, MPI_DOUBLE, 1, 0, MPI_COMM_WORLD); // to rank 1, tag 0
} else if (rank == 1) {
MPI_Recv(&x, 1, MPI_DOUBLE, 0, 0, MPI_COMM_WORLD, MPI_STATUS_IGNORE);
}
MPI_Finalize();
return 0;
}
The same program runs on every node, and each copy decides what to do from its rank. Beyond point-to-point messages, MPI provides collective operations that involve every rank: broadcast a value from one to all, gather results, or all-reduce, where every rank contributes a number and every rank gets the sum. All-reduce is also the core of training neural networks on many GPUs: after each step, the gradients computed on each GPU are summed across all of them.
Message passing is harder to program than shared memory, because the programmer has to decide where every piece of data lives and when it moves. It has one big advantage: there is no hidden traffic. Nothing like false sharing can happen by accident, and the cost of communication is visible in the code.
What a message costs
The cost of sending a message is roughly a fixed latency plus the size divided by the bandwidth. Latency is what it costs to send one byte: software on both sides, the network interface, the wires and switches in between. Here is the round trip for a tiny message at different distances, measured from the M2 Ultra:
| Path | Round trip for a small message |
|---|---|
| a cache line between two cores (shared memory) | 66–400 ns |
| TCP over loopback, same machine, macOS | 19–20 µs |
| TCP over loopback, Linux VM on the same Mac | 36–37 µs |
| TCP handshake with example.com's nearest server | 6.4–6.9 ms |
The loopback numbers come from a C program that bounces 1 byte back and forth 20,000 times between two threads through a TCP socket on 127.0.0.1, with Nagle's algorithm off. No wire is involved, and it's still 200 times slower than moving a cache line. The time goes into the kernel: two system calls per direction, the TCP/IP stack, and waking the sleeping receiver thread. The last line is the TCP handshake time for HTTPS requests to example.com, as reported by curl (connect time minus name lookup): one real round trip over the internet, about 300 times longer again.
Bandwidth is the other half. The same loopback program with 1 MiB messages took 211–224 µs per round trip, about 9.4–9.9 GB/s in each direction. With a one-way latency of about 10 µs and 10 GB/s, a 1 MiB message is dominated by the size, and an 8-byte message entirely by the latency. The size where both count equally is latency × bandwidth, about 100 KB here. Try it:
Try it: Press Next line to run the highlighted line, or Play to watch; the boxes on the right are the variables, and the ones that just changed light up. Edit the code to try your own changes.
1// Time to send one message: latency + size / bandwidth.2long latency_ns = 10000; // 10 us one way (loopback TCP on the M2)3long bytes_per_ns = 10; // 10 GB/s45int main() {6 long size = 8;7 while (size <= 4194304) {8 long t = latency_ns + size / bytes_per_ns;9 long mb_s = size * 1000 / t; // effective MB/s10 printf("%8ld bytes: %8ld ns, %5ld MB/s (%ld%% of peak)\n",11 size, t, mb_s, mb_s / (bytes_per_ns * 10));12 size = size * 8;13 }14 return 0;15}
At 10 µs, a 4 KB message gets 3% of the peak bandwidth and a 256 KB one 72%. Change latency_ns to 1000, closer to what a supercomputer network offers, and 32 KB messages already reach 76%. Small messages are all latency, which is why parallel programs try to send fewer, bigger messages, and why HPC networks fight for every microsecond.
They do it by getting the kernel out of the way. Networks like InfiniBand and HPE's Slingshot use RDMA (remote direct memory access): a process registers a buffer with the network card once, and after that the card reads and writes it directly, with no system call and no copy, and can write straight into a buffer on the remote node. Small-message latency between two nodes then falls to the order of a microsecond or two, most of it in the network cards and switches.
Interconnection networks
With thousands of nodes, the wiring pattern (the topology) matters as much as the speed of each link. Two measures describe it:
- the diameter: the largest number of hops between any two nodes, which bounds the worst-case latency;
- the bisection bandwidth: the number of links that must be cut to split the machine into two halves, which bounds how much data can cross the middle when every node talks to a far-away partner.
For 64 nodes:
| Topology | Links | Diameter | Bisection (links) |
|---|---|---|---|
| ring | 64 | 32 | 2 |
| 8 × 8 mesh | 112 | 14 | 8 |
| 8 × 8 torus (mesh with wraparound) | 128 | 8 | 16 |
| 6-dimensional hypercube | 192 | 6 | 32 |
| fully connected | 2,016 | 1 | 1,024 |
The ring is cheap and hopeless: half the traffic has to squeeze through 2 links. Full connection is ideal and impossible beyond a handful of nodes. Real machines sit in between. IBM's Blue Gene machines used 3D and then 5D tori; Japan's Fugaku uses a six-dimensional mesh/torus called Tofu. Most clusters use a fat tree built from switches, where links get fatter (more of them in parallel) toward the root, so that bisection bandwidth doesn't shrink at the top. HPE's Slingshot uses a dragonfly: groups of switches fully connected inside, with every group directly linked to every other group, so any two nodes are only a few switch hops apart.
How a message crosses the network matters too. Store-and-forward switching receives a whole packet at each switch before sending it on, so latency grows with the number of hops times the packet size. Cut-through and wormhole switching start forwarding a packet as soon as its header has arrived and the route is decided, so the packet streams through several switches at once and the extra cost per hop is small. High-performance networks all use the second approach.
From MPPs to clusters
Early large multicomputers were MPPs (massively parallel processors): thousands of processors with a custom network, custom packaging and often custom chips, from vendors like Cray, Intel, Thinking Machines and IBM. They were expensive and tied to one vendor.
In 1994, a project at NASA called Beowulf showed a cheaper way: a cluster of ordinary PCs running Linux, connected by ordinary Ethernet, programmed with message passing. Performance per dollar was hard to beat, and clusters of commodity servers took over most of high-performance computing.
The line between the two has blurred since. Today's top systems are built from nodes that use standard server processors and GPUs, run Linux, like the rest of the list (see the Unix and Windows chapter), but are packaged by a few vendors with liquid cooling and a specialized network. They are clusters in architecture and MPPs in engineering.
Warehouse-scale computers
The largest clusters of all don't run scientific simulations. A company like Google, Amazon, Microsoft or Meta runs data centers with tens of thousands of servers each, and treats a whole building as one computer running a few huge services. Google's engineers called them warehouse-scale computers in a 2009 book on the subject.
Their design differs from supercomputers in three ways:
- Failure is normal. With that many disks, memory modules and power supplies, something breaks every day. Software is designed to tolerate it: data is replicated on several machines, and frameworks like MapReduce (published by Google in 2004) rerun the work of a failed machine elsewhere. A supercomputer job, by contrast, usually stops when a node dies and restarts from its last checkpoint.
- Throughput over latency. A search engine handles millions of independent requests, each of which needs only a small part of the machine. That parallelism is easy; the hard part is the tail: when one request fans out to hundreds of servers, the slowest one decides the response time.
- Energy and cost. The building's power and cooling matter as much as the servers. The standard measure, PUE (power usage effectiveness), divides the whole facility's power by the power that reaches the computers; the best data centers get close to 1, meaning almost no energy is spent on anything but computing.
Cloud computing sells slices of these machines. It largely replaced the earlier idea of grid computing, which tried to federate the separate clusters of many institutions into one shared resource.
Supercomputers and the TOP500
Since June 1993, the TOP500 list has ranked the world's fastest supercomputers twice a year, in June and November. The ranking uses one benchmark, HPL (High-Performance LINPACK), which solves a huge dense system of linear equations in 64-bit floating point. Each entry lists Rmax, the speed achieved on HPL, and Rpeak, the theoretical maximum of the hardware.
In June 2022, Frontier, at Oak Ridge National Laboratory in the United States, became the first system to pass an exaflop (10¹⁸ floating-point operations per second) on HPL, with an Rmax of 1.102 exaflops. It's built from HPE Cray EX nodes, each combining an AMD EPYC processor with AMD Instinct MI250X GPUs, linked by Slingshot-11. El Capitan, at Lawrence Livermore, took first place in November 2024.
The top five of the June 2026 list, the 67th edition:
| Rank | System | Site | Rmax (exaflops) | Power (MW) | Processors |
|---|---|---|---|---|---|
| 1 | LineShine | NSC Shenzhen, China | 2.198 | 42.2 | 304-core LX2 CPUs only |
| 2 | El Capitan | LLNL, USA | 1.809 | 29.7 | AMD EPYC + MI300A |
| 3 | Frontier | Oak Ridge, USA | 1.353 | 24.6 | AMD EPYC + MI250X GPUs |
| 4 | Aurora | Argonne, USA | 1.012 | 38.7 | Intel Xeon Max + Intel GPUs |
| 5 | JUPITER Booster | Jülich, Germany | 1.000 | 15.8 | NVIDIA GH200 |
Some things to read in it:
- Five exascale systems. Frontier was alone for two years; now there are five, and Frontier's own score grew from 1.102 to 1.353 as the system was tuned.
- Rmax versus Rpeak. Frontier reaches 66% of its peak and El Capitan 64%; LineShine 80%. Even the most regular computation there is doesn't keep all units busy, because data has to be moved and exchanged.
- Power. An exaflop costs tens of megawatts, the consumption of a small town. El Capitan does about 61 billion operations per joule, Frontier 55, LineShine 52. Energy, not transistor count, now limits how big these machines get, for the same reason clock speeds stopped rising.
- GPUs. Four of the five get most of their speed from GPUs or GPU-like accelerators, for the reasons in the SIMD and GPUs chapter: more arithmetic per watt on regular, data-parallel work. The new number one is the exception, built from CPUs with 304 cores each, as Fugaku was with its ARM A64FX processors when it led the list in 2020 and 2021.
HPL is a flattering benchmark: dense matrix arithmetic, with a lot of computation per byte moved. Many real applications are limited by memory bandwidth and communication instead. The companion HPCG benchmark, dominated by sparse, memory-bound operations, gives LineShine 22.00 petaflops on the same list, about 1% of its HPL score.
Running a program on 10 million cores
Amdahl's law is brutal at this scale: with 10 million cores, even a serial fraction of one in a million caps the speedup at about 900,000. Supercomputers are still useful because problems grow with the machine. A climate model on a bigger computer uses a finer grid, not the same grid faster. This weak scaling keeps each node's share of the work roughly constant.
The usual structure is domain decomposition: the simulated space is cut into blocks, one per process. Each process updates its own block, then exchanges a thin layer of boundary cells, the halo, with its neighbours. Computation grows with a block's volume and communication with its surface, so bigger blocks mean relatively less communication.
Nodes are allocated by a batch scheduler such as Slurm: a job asks for a number of nodes and a time limit, waits in a queue, and gets its nodes to itself for the duration. Because the processes of a parallel job wait on each other at every exchange, they must all run at once; the scheduler never time-slices one job's nodes with another's.
People have long tried to make a multicomputer look like shared memory in software. Distributed shared memory systems, starting in the 1980s, used the virtual memory hardware to fetch a page from another node on a page fault. It works, but false sharing at page granularity and the cost of each fault made it slow. Today's compromise is the PGAS model (partitioned global address space), in languages like UPC and Chapel: one address space in the program, but every piece of data has a visible home, so the programmer knows which accesses are remote.
Takeaways
- A multicomputer or cluster is many computers with private memories cooperating by message passing, usually through MPI: send, receive and collectives like all-reduce.
- Message cost ≈ latency + size / bandwidth. Measured from the M2 Ultra: 66–400 ns to move a cache line between cores, 19–20 µs for a loopback TCP round trip, 6.4–6.9 ms to a web server. Loopback moved about 9.4–9.9 GB/s with 1 MiB messages.
- HPC networks bypass the kernel with RDMA to bring node-to-node latency to about a microsecond.
- Topologies trade cost against diameter and bisection bandwidth: rings, meshes, tori, fat trees, dragonflies. Cut-through switching keeps per-hop cost small.
- Clusters of commodity servers (Beowulf, 1994) replaced most custom MPPs. Warehouse-scale computers assume failure is normal and optimize throughput, tail latency and energy.
- The TOP500 ranks systems on HPL. Frontier passed 1 exaflop in June 2022; the June 2026 list has five exascale systems, led by the CPU-only LineShine at 2.198 exaflops, and four of the five rely on GPUs. Each draws 16 to 42 MW.