Skip to content

Niveau 6 · Chapitre 6.12

Multiprocesseurs à mémoire partagée et NUMA

Beaucoup de processeurs, un seul espace d’adressage : multiprocesseurs et multiordinateurs, machines UMA à bus et pourquoi elles plafonnent, crossbars et réseaux de commutation, NUMA et placement au premier accès, cohérence par répertoire, modèles de cohérence mémoire, et les deux puces du M2 Ultra.

Le chapitre sur les multicœurs a mis plusieurs cœurs sur une puce, autour d’une seule mémoire. Les serveurs vont plus loin : deux, quatre ou huit puces de processeur sur une même carte, parfois plusieurs puces (dies) dans chaque boîtier, qui partagent toutes un même espace d’adressage. N’importe quel cœur peut charger n’importe quelle adresse, et le matériel fait comme s’il n’y avait qu’une grande mémoire.

Cette illusion devient plus difficile à tenir à mesure que la machine grandit. La mémoire ne peut pas être également proche de centaines de cœurs, un bus partagé unique ne peut pas porter le trafic de tous, et diffuser chaque défaut de cache à tous les caches cesse de fonctionner. Ce chapitre décrit comment on construit les grosses machines à mémoire partagée : le réseau d’interconnexion, la mémoire NUMA, la cohérence par répertoire et les règles de cohérence mémoire qu’elles offrent au logiciel.

Multiprocesseurs et multiordinateurs

Les ordinateurs parallèles se rangent en deux familles selon la façon dont leurs processeurs communiquent :

MultiprocesseurMultiordinateur
Mémoireun espace d’adressage partagéchaque nœud a sa propre mémoire privée
Communicationchargements et rangements ordinairesmessages explicites sur un réseau
Programmationthreads, verrous, atomiquespassage de messages (sockets, MPI)
Taillejusqu’à quelques centaines de cœursjusqu’à des millions de cœurs
Exempleun portable, un serveur bi-socketune grappe, un supercalculateur

Un multiprocesseur est plus facile à programmer : les structures de données sont simplement partagées, comme dans le chapitre sur les threads. Un multiordinateur est plus facile à construire en grand : chaque nœud est un ordinateur ordinaire et le réseau ne transporte que des messages. Le chapitre sur les grappes traite des multiordinateurs. Ici, on reste en mémoire partagée.

UMA : le multiprocesseur à bus

Le multiprocesseur le plus simple place plusieurs processeurs et une mémoire sur un bus partagé. Chaque processeur atteint chaque adresse dans le même temps : on parle d’UMA (uniform memory access, accès mémoire uniforme), ou de SMP (symmetric multiprocessor) puisqu’aucun processeur n’est privilégié.

Les caches rendent la chose possible. Chaque processeur garde ses données de travail dans son cache et n’utilise le bus qu’en cas de défaut, et les caches restent cohérents par espionnage (snooping) : chaque cache observe chaque transaction du bus et invalide ou fournit sa copie comme l’exige le protocole MESI. Un bus diffuse naturellement à tous, donc l’espionnage est simple à construire.

Le bus est aussi la limite. Chaque défaut, chaque réécriture et chaque invalidation de chaque processeur passe par les mêmes fils, un à la fois. Ajouter des processeurs ajoute du trafic mais pas de débit, et au-delà de quelques dizaines de processeurs le bus sature. Les SMP à bus des années 1990 s’arrêtaient à peu près à cette taille. Aujourd’hui, même une seule puce a dépassé le bus : les cœurs d’un gros processeur de serveur sont reliés par un anneau ou un maillage, et une seule puce multicœur peut contenir plus de cœurs que la plus grosse machine à bus n’en a jamais eu.

Crossbars et réseaux de commutation

L’alternative au bus est le commutateur. Un crossbar (commutateur à barres croisées) relie n processeurs à n modules mémoire par une grille de n² points de croisement : n’importe quel processeur peut joindre n’importe quel module libre, et les n connexions peuvent être actives en même temps. Le problème, c’est le n² : 64 points de croisement pour 8 processeurs, mais plus d’un million pour 1 024.

Un réseau multiétage sacrifie un peu de ce parallélisme contre beaucoup moins de commutateurs. Le réseau oméga enchaîne log₂ n étages de petits commutateurs 2 × 2, n / 2 par étage. Pour 8 entrées, cela fait 3 étages de 4 commutateurs : 12 commutateurs au lieu de 64 points de croisement. Pour 1 024 entrées, 10 étages de 512 : 5 120 commutateurs au lieu de 1 048 576. Chaque commutateur aiguille une requête d’après un bit de l’adresse de destination, sans contrôle central. Le prix est le blocage : deux requêtes peuvent avoir besoin du même lien interne au même moment, même si leurs destinations diffèrent, et l’une doit attendre.

Les machines actuelles n’utilisent ni un unique crossbar ni un réseau oméga d’école, mais le compromis est le même. Dans une puce, cœurs, tranches de cache et contrôleurs mémoire sont placés sur un anneau ou un maillage 2D, chaque saut coûtant un ou deux cycles. Entre puces, des liens point à point (l’UPI d’Intel, l’Infinity Fabric d’AMD) relient directement les sockets. Aucun lien n’est partagé par tous, donc le débit croît avec le nombre de processeurs. Mais la distance à la mémoire n’est plus la même pour tous.

NUMA

Dès que la mémoire est répartie dans la machine, chaque morceau est proche de certains processeurs et loin des autres. Chaque puce de processeur a ses propres contrôleurs mémoire et sa propre DRAM. Un chargement vers cette mémoire locale y va directement. Un chargement vers la mémoire rattachée à un autre socket doit traverser l’interconnexion, aller et retour. Les deux fonctionnent (c’est toujours un seul espace d’adressage, accédé par des chargements ordinaires), mais ils ne prennent pas le même temps. C’est le NUMA (non-uniform memory access, accès mémoire non uniforme).

Un processeur avec sa mémoire locale forme un nœud NUMA. Sous Linux, numactl --hardware affiche la disposition, dont une table des distances tirée des tables ACPI du micrologiciel, où l’accès local vaut 10 par définition et où les distances distantes sont relatives : un serveur bi-socket typique indique environ 20 pour l’autre socket, c’est-à-dire une mémoire distante à peu près deux fois plus loin. La table n’est qu’une estimation ; les vraies latences, ce sont les mesures qui les donnent.

Presque toutes les grosses machines sont NUMA aujourd’hui, et forcément NUMA à caches cohérents (CC-NUMA) : les caches restent cohérents partout, y compris à travers l’interconnexion. Des machines sans cohérence matérielle (NC-NUMA) existent, mais elles sont difficiles à programmer.

Le NUMA apparaît même à l’intérieur d’un boîtier. Les processeurs EPYC d’AMD construisent un socket à partir d’une douzaine de petites puces de cœurs (chiplets) ou plus, autour d’une puce d’E/S centrale qui porte les contrôleurs mémoire, et peuvent présenter ce socket comme un, deux ou quatre nœuds NUMA (le réglage NPS). Le sub-NUMA clustering d’Intel découpe de même un gros Xeon en plusieurs nœuds. Et CXL (Compute Express Link), un protocole cohérent avec les caches qui passe par PCI Express, permet d’ajouter de la mémoire sur une carte ; Linux la présente comme un nœud NUMA doté de mémoire mais sans processeur.

Placer la mémoire

Sur une machine NUMA, la performance dépend de l’endroit où vit chaque page. La politique par défaut de Linux est le premier accès (first touch) : une page physique est allouée sur le nœud du processeur qui y écrit le premier, au moment du défaut de page, pas là où malloc a été appelé. D’où un piège classique : un programme qui initialise un grand tableau dans son thread principal, puis le traite avec 64 threads, a tout placé sur un seul nœud. Tous les autres nœuds lisent à distance, et les contrôleurs mémoire de ce nœud deviennent le goulet d’étranglement. Le remède : initialiser les données en parallèle, avec les threads qui s’en serviront.

Linux permet aussi de choisir explicitement : numactl --cpunodebind=0 --membind=0 ./prog exécute un programme sur le nœud 0 avec de la mémoire du nœud 0 uniquement, --interleave=all répartit les pages à tour de rôle sur tous les nœuds (bien pour des données que tout le monde utilise), et les appels système mbind et set_mempolicy font la même chose depuis le code. L’ordonnanceur essaie de garder un thread sur le nœud où se trouve sa mémoire, et peut migrer les pages vers les threads qui les utilisent.

La cohérence par répertoire

L’espionnage suppose que tout le monde voit chaque requête. Sur une machine NUMA de dizaines de sockets, diffuser chaque défaut à chaque cache noierait l’interconnexion. La réponse est un répertoire (directory) : pour chaque ligne de mémoire, la liste des caches qui la détiennent.

Chaque ligne a un nœud d’origine (home node), celui dont la mémoire la contient, et c’est lui qui tient son entrée de répertoire : un état (non cachée, partagée ou modifiée) et un vecteur de bits, un bit par nœud. Un défaut en lecture va au nœud d’origine, pas à tout le monde :

  • Si la ligne est non cachée ou partagée, le nœud d’origine envoie les données et positionne le bit du demandeur.
  • Si un nœud la détient modifiée, le nœud d’origine lui transmet la requête ; ce nœud envoie les données et rétrograde sa copie.

Une écriture passe aussi par le nœud d’origine. Il envoie une invalidation exactement aux nœuds dont le bit est positionné, recueille leurs acquittements, puis accorde la propriété. Le trafic est point à point et proportionnel au nombre de détenteurs réels, qui est le plus souvent zéro ou un.

Le répertoire coûte de la mémoire. Un vecteur complet pour 64 nœuds prend 64 bits par ligne : 12,5 % de plus pour une ligne de 64 octets. Les conceptions réelles le réduisent avec des vecteurs plus grossiers (un bit par groupe de nœuds) ou en ne gardant d’entrées que pour les lignes effectivement cachées quelque part, dans un cache de répertoire. Le prototype DASH de Stanford au début des années 1990, puis les serveurs Origin de SGI plus tard dans la décennie, ont rendu pratique le CC-NUMA à répertoire, et la même idée sert désormais à toutes les échelles : les grosses puces tiennent un répertoire ou un filtre d’espionnage à côté de chaque tranche de cache, et les serveurs multisockets suivent quel socket détient chaque ligne.

Une conception plus radicale, COMA (cache-only memory architecture), transformait toute la mémoire principale en un gigantesque cache, pour que les données migrent vers le nœud qui s’en sert. Elle a été essayée commercialement vers 1990 et n’a pas survécu ; la migration de pages par le système d’exploitation en obtient une partie du bénéfice, plus simplement.

Ce que voient les autres processeurs

La cohérence de cache concerne un seul emplacement : tous les processeurs s’accordent sur l’ordre des écritures dans x. Elle ne dit rien de l’ordre dans lequel ils voient les écritures dans x et dans y. C’est le rôle du modèle de cohérence mémoire (memory consistency model), et un gros multiprocesseur a toutes les raisons d’en vouloir un faible : ses tampons d’écriture sont profonds, et les écritures vers des lignes dont les nœuds d’origine diffèrent se terminent naturellement à des moments différents.

Les modèles possibles, du plus fort au plus faible :

  • Cohérence stricte : chaque lecture renvoie l’écriture la plus récente en temps réel. Aucune machine réelle ne peut l’offrir : il faudrait décider instantanément, dans toute la machine, laquelle est « la plus récente ».
  • Cohérence séquentielle (Leslie Lamport, 1979) : le résultat est le même que si toutes les opérations de tous les processeurs s’exécutaient dans un unique entrelacement, en respectant l’ordre du programme de chacun. C’est ce que les programmeurs supposent naturellement, et c’est trop coûteux à fournir à pleine vitesse.
  • Cohérence processeur / TSO : les écritures de chaque processeur sont vues par tous dans l’ordre où il les a émises, mais un chargement peut doubler un rangement antérieur du même processeur. C’est le x86.
  • Ordre faible et cohérence à la libération (release consistency) : presque tous les réordonnancements sont permis, et le programme marque les points où l’ordre compte : barrières, ou opérations d’acquisition et de libération autour des sections critiques. C’est l’ARM64 et le RISC-V.

La vue d’ensemble de l’ISA exécute les tests décisifs (litmus tests) sur le M2 : le réordonnancement du test store buffering est apparu dans 91 % des exécutions en natif et 97 % dans le mode x86 de Rosetta, et a disparu avec une barrière. La règle pratique est la même sur toute machine : faire communiquer les threads par des atomiques et des verrous, qui contiennent les bonnes barrières, et jamais par des variables ordinaires.

Le M2 Ultra : deux puces, une mémoire

Le M2 Ultra sur lequel ce site a été mesuré est lui-même un petit multiprocesseur. Ce sont deux puces M2 Max réunies par la connexion UltraFusion d’Apple, un interposeur en silicium qui, selon Apple, porte plus de 10 000 signaux et 2,5 To/s entre les deux puces. Chaque puce a 8 cœurs performance, 4 cœurs efficacité et ses propres contrôleurs mémoire ; ensemble, elles atteignent 800 Go/s de débit mémoire, le double des 400 Go/s d’un M2 Max.

Physiquement, chaque puce est donc plus proche de sa moitié de la mémoire. Mais le système n’en montre rien :

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

Un boîtier, 24 cœurs, et aucune interface NUMA dans macOS. La mesure va dans le même sens. Un parcours aléatoire de pointeurs sur 512 Mo, lancé depuis 60 threads fraîchement créés et placés là où l’ordonnanceur l’a voulu, a donné 131 à 137 ns par chargement à chaque fois, sans aucun second groupe d’accès plus lents. Soit l’ordonnanceur a gardé chaque thread près de ses données, ce qui est peu probable sur 60 essais, soit les adresses sont réparties finement entre les contrôleurs des deux puces, de sorte que chaque cœur voit le même mélange. Apple ne documente pas lequel ; elle présente la puce comme une mémoire uniforme, et pour le logiciel elle se comporte comme telle.

Le trafic de cohérence, c’est une autre affaire. Une ligne de cache qui rebondit entre deux cœurs a pris de 66 à environ 400 ns par aller-retour dans le chapitre sur les multicœurs, selon l’endroit où tournaient les deux threads. Même une machine à mémoire uniforme n’est pas uniforme pour des données partagées et écrites.

À retenir

  • Un multiprocesseur partage un espace d’adressage et communique par chargements et rangements ; un multiordinateur a des mémoires privées et communique par messages.
  • Les machines UMA placent tous les processeurs sur un bus, avec des caches espions ; le bus sature au-delà de quelques dizaines de processeurs.
  • Un crossbar demande n² points de croisement ; un réseau oméga (n/2)·log₂ n commutateurs (5 120 au lieu de 1 048 576 pour 1 024 ports), mais il peut bloquer. Les puces actuelles utilisent des anneaux et des maillages ; les sockets, des liens point à point.
  • En NUMA, chaque nœud a sa mémoire locale ; la mémoire distante est accessible par des chargements ordinaires, mais plus lentement. Linux alloue au premier accès : initialisez les données avec les threads qui s’en serviront ; numactl contrôle le placement.
  • La cohérence par répertoire note les détenteurs de chaque ligne à son nœud d’origine et n’envoie les invalidations qu’à eux, au lieu de diffuser.
  • Les modèles de cohérence vont de la cohérence séquentielle au TSO (x86) et à l’ordre faible / à la libération (ARM64, RISC-V) ; barrières et atomiques rétablissent l’ordre.
  • Le M2 Ultra est fait de deux puces réunies par UltraFusion, présentées comme un seul boîtier à mémoire uniforme : 131 à 137 ns de latence DRAM depuis chacun de 60 threads.

Dans ce niveau

  1. 6.1Le cycle fetch–decode–execute
  2. 6.2Chemin de données et bus
  3. 6.3Unité de contrôle et microcode
  4. 6.4Une machine complète : la Mic-1 exécutant IJVM
  5. 6.5Pipeline et aléas
  6. 6.6Caches et hiérarchie mémoire
  7. 6.7Prédiction de branchement
  8. 6.8Exécution dans le désordre, renommage de registres et spéculation
  9. 6.9Cœurs réels : x86, ARM et AVR comparés
  10. 6.10SIMD, GPU et coprocesseurs
  11. 6.11Multicœurs, multithreading et cohérence de cache
  12. 6.12Multiprocesseurs à mémoire partagée et NUMA
  13. 6.13Grappes, passage de messages et supercalculateurs