Skip to content

Level 6 · Chapter 6.12

Shared-memory multiprocessors and NUMA

Many processors, one address space: multiprocessors versus multicomputers, bus-based UMA machines and why they stop scaling, crossbars and switching networks, NUMA and first-touch placement, directory-based coherence, consistency models, and the M2 Ultra's two dies.

The multicore chapter put several cores on one chip, sharing one memory. Servers go further: two, four or eight processor chips on one board, sometimes several dies in each package, all sharing one address space. Any core can load from any address, and the hardware makes it look like one big memory.

That illusion gets harder to keep as machines grow. Memory can't be equally close to hundreds of cores, a single shared bus can't carry everyone's traffic, and broadcasting every cache miss to every cache stops working. This chapter is about how large shared-memory machines are built: the interconnect, NUMA memory, directory-based coherence and the consistency rules they give software.

Multiprocessors and multicomputers

Parallel computers split into two families according to how their processors communicate:

MultiprocessorMulticomputer
Memoryone shared address spaceeach node has its own private memory
Communicationordinary loads and storesexplicit messages over a network
Programmingthreads, locks, atomicsmessage passing (sockets, MPI)
Sizeup to a few hundred coresup to millions of cores
Examplea laptop, a 2-socket servera cluster, a supercomputer

A multiprocessor is easier to program: data structures are simply shared, as in the threads chapter. A multicomputer is easier to build big: each node is an ordinary computer and the network only carries messages. The clusters chapter covers multicomputers. Here we stay with shared memory.

UMA: the bus-based multiprocessor

The simplest multiprocessor puts several CPUs and one memory on a shared bus. Every CPU reaches every address in the same time, so it's called UMA (uniform memory access), or SMP (symmetric multiprocessor) because no processor is special.

Caches make it work. Each CPU keeps its working data in its cache and only uses the bus on a miss, and the caches stay coherent by snooping: every cache watches every bus transaction and invalidates or supplies its copy as the MESI protocol dictates. A bus is a natural broadcast medium, so snooping is cheap to build.

The bus is also the limit. Every miss, every write-back and every invalidation from every CPU goes over the same wires, one at a time. Adding CPUs adds traffic but no bandwidth, and past a few dozen processors the bus saturates. Bus-based SMPs of the 1990s stopped at about that size. Today, even a single chip has outgrown a bus: the cores of a large server processor are connected by a ring or a mesh, and one multicore chip can hold more cores than the biggest bus-based machine ever had.

Crossbars and switching networks

The alternative to a bus is a switch. A crossbar connects n CPUs to n memory modules through a grid of n² crosspoints: any CPU can reach any free memory module, and all n connections can be active at once. The problem is the n²: 64 crosspoints for 8 CPUs, but over a million for 1,024.

A multistage network trades some of that parallelism for far fewer switches. The omega network chains log₂ n stages of small 2 × 2 switches, n / 2 per stage. For 8 inputs that's 3 stages of 4 switches: 12 switches instead of 64 crosspoints. For 1,024 inputs, 10 stages of 512: 5,120 switches instead of 1,048,576. Each switch routes a request with one bit of the destination address, so routing needs no central control. The price is blocking: two requests may need the same internal link at the same time, even when their destinations differ, and one has to wait.

Current machines use neither a single crossbar nor a pure omega network, but the trade-off is the same one. Inside a chip, cores, cache slices and memory controllers sit on a ring or a 2D mesh, each hop costing a cycle or two. Between chips, point-to-point links (Intel's UPI, AMD's Infinity Fabric) connect sockets directly. No link is shared by everyone, so bandwidth grows with the number of processors. But the distance to memory is no longer the same for all of them.

NUMA

Once memory is spread over the machine, each part of it is close to some processors and far from others. Each processor chip has its own memory controllers and its own DRAM. A load to that local memory goes straight to it. A load to memory attached to another socket has to cross the interconnect and back. Both work (it's still one address space, accessed with ordinary loads), but they take different times. That's NUMA (non-uniform memory access).

A processor with its local memory is a NUMA node. Linux shows the layout with numactl --hardware, including a distance table from the firmware's ACPI tables, where local access is defined as 10 and remote distances are relative to it: a typical two-socket server reports around 20 for the other socket, meaning remote memory is roughly twice as far. The table is only an estimate; the real latencies are for benchmarks to measure.

Almost every large machine is NUMA today, and it has to be cache-coherent NUMA (CC-NUMA): caches everywhere stay coherent even across the interconnect. Machines without hardware coherence (NC-NUMA) exist but are hard to program.

NUMA even appears within one package. AMD's EPYC processors build a socket from a dozen or more core chiplets around a central I/O die that holds the memory controllers, and can present that socket as one, two or four NUMA nodes (the NPS setting). Intel's sub-NUMA clustering splits a large Xeon into several nodes the same way. And CXL (Compute Express Link), a cache-coherent protocol over PCI Express, lets extra memory be attached on a card; Linux shows it as a NUMA node with memory but no CPUs.

Placing memory

On a NUMA machine, performance depends on where each page lives. Linux's default policy is first touch: a physical page is allocated on the node of the CPU that first writes to it, when the page fault happens, not where malloc was called. That has a classic pitfall: a program that initializes a large array in its main thread, then processes it with 64 threads, has placed all of it on one node. Every other node reads remotely, and that one node's memory controllers become the bottleneck. The fix is to initialize data in parallel, with the same threads that will use it.

Linux also lets a program choose explicitly: numactl --cpunodebind=0 --membind=0 ./prog runs a program on node 0 with memory from node 0 only, --interleave=all spreads pages round-robin over all nodes (good for data everyone uses), and the mbind and set_mempolicy system calls do the same from code. The scheduler tries to keep a thread on the node where its memory is, and can migrate pages toward the threads that use them.

Directory-based coherence

Snooping relies on everyone seeing every request. On a NUMA machine with dozens of sockets, broadcasting each miss to every cache would flood the interconnect. The answer is a directory: for each line of memory, a record of which caches hold it.

Each line has a home node, the one whose memory contains it, and the home keeps its directory entry: a state (uncached, shared, or modified) and a bit vector with one bit per node. A read miss goes to the home, not to everyone:

  • If the line is uncached or shared, the home sends the data and sets the requester's bit.
  • If one node holds it modified, the home forwards the request to that node, which sends the data and downgrades its copy.

A write goes to the home too. The home sends an invalidation to exactly the nodes whose bits are set, collects their acknowledgements, then grants ownership. The traffic is point-to-point and proportional to the number of actual sharers, which is usually zero or one.

The directory has a cost in memory. A full bit vector for 64 nodes takes 64 bits per line: 12.5% extra for a 64-byte line. Real designs shrink it with coarser vectors (one bit per group of nodes) or by keeping entries only for lines that are actually cached somewhere, in a directory cache. Stanford's DASH prototype in the early 1990s and SGI's Origin servers later in the decade made directory-based CC-NUMA practical, and the same idea now works at every scale: large chips keep a directory or snoop filter next to each cache slice, and multi-socket servers track which socket holds each line.

A more radical design, COMA (cache-only memory architecture), turned all of main memory into a giant cache, so data migrated to whichever node used it. It was tried commercially around 1990 and didn't survive; page migration by the operating system gets part of the benefit more simply.

What other processors see

Coherence is about a single location: all processors agree on the order of writes to x. It says nothing about the order in which they see writes to x and y. That is the memory consistency model, and a large multiprocessor has every reason to want a weak one: its store buffers are deep, and writes to lines homed on different nodes naturally complete at different times.

The possible models, from strongest to weakest:

  • Strict consistency: every read returns the most recent write in real time. No real machine can offer it: "most recent" would have to be decided instantly across the whole machine.
  • Sequential consistency (Leslie Lamport, 1979): the result is as if all operations of all processors were executed in some single interleaving, keeping each processor's program order. It's what programmers naturally assume, and too expensive to provide at full speed.
  • Processor consistency / TSO: each processor's writes are seen by everyone in the order it issued them, but a processor's load may overtake its own earlier store. That's x86.
  • Weak ordering and release consistency: almost any reordering is allowed, and the program marks the points where order matters: fences, or acquire and release operations around critical sections. That's ARM64 and RISC-V.

The ISA overview runs the litmus tests on the M2: the store-buffering reordering showed up in 91% of runs natively and 97% under Rosetta's x86 mode, and disappeared with a fence. The practical rule is the same on any machine: communicate between threads through atomics and locks, which include the right fences, and never through plain variables.

The M2 Ultra: two dies, one memory

The M2 Ultra that this site was measured on is itself a small multiprocessor. It's two M2 Max dies joined by Apple's UltraFusion connection, a silicon interposer that Apple describes as carrying more than 10,000 signals and 2.5 TB/s between the dies. Each die has 8 performance cores, 4 efficiency cores and its own memory controllers; together they reach 800 GB/s of memory bandwidth, twice an M2 Max's 400 GB/s.

Physically, then, each die is closer to its own half of the memory. But the operating system shows none of it:

$ sysctl hw.packages machdep.cpu.cores_per_package hw.ncpu
hw.packages: 1
machdep.cpu.cores_per_package: 24
hw.ncpu: 24

One package, 24 cores, and no NUMA interface in macOS. Measurement agrees with that picture. A random pointer chase through 512 MB, run from 60 freshly created threads placed wherever the scheduler chose, gave 131–137 ns per load every time, with no second group of slower accesses. Either the scheduler kept every thread near its data, which is unlikely over 60 tries, or addresses are spread finely across both dies' controllers so that every core sees the same mix. Apple doesn't document which; it presents the chip as uniform memory, and for software it behaves as such.

Coherence traffic is another matter. A cache line bouncing between two cores took anywhere from 66 to about 400 ns per round trip in the multicore chapter, depending on where the two threads happened to run. Even a machine with uniform memory isn't uniform for shared, written data.

Takeaways

  • A multiprocessor shares one address space and communicates by loads and stores; a multicomputer has private memories and communicates by messages.
  • UMA machines put all CPUs on one bus with snooping caches; the bus saturates beyond a few dozen processors.
  • A crossbar needs n² crosspoints; an omega network needs (n/2)·log₂ n switches (5,120 instead of 1,048,576 for 1,024 ports) but can block. Chips today use rings and meshes; sockets use point-to-point links.
  • In NUMA, each node has local memory; remote memory is reachable by ordinary loads but slower. Linux allocates on first touch, so initialize data with the threads that will use it; numactl controls placement.
  • Directory-based coherence records the sharers of each line at its home node and sends invalidations only to them, instead of broadcasting.
  • Consistency models range from sequential consistency to TSO (x86) and weak/release ordering (ARM64, RISC-V); fences and atomics restore order.
  • The M2 Ultra is two dies joined by UltraFusion, presented as one package with uniform memory: 131–137 ns DRAM latency from any of 60 threads.

In this level

  1. 6.1The fetch–decode–execute cycle
  2. 6.2Datapath, internal and system buses
  3. 6.3Control units and microcode
  4. 6.4A complete machine: the Mic-1 running IJVM
  5. 6.5Pipelining and hazards
  6. 6.6Caches and the memory hierarchy
  7. 6.7Branch prediction
  8. 6.8Out-of-order execution, register renaming and speculation
  9. 6.9Real cores: x86, ARM and AVR compared
  10. 6.10SIMD, GPUs and coprocessors
  11. 6.11Multicore, multithreading and cache coherence
  12. 6.12Shared-memory multiprocessors and NUMA
  13. 6.13Clusters, message passing and supercomputers