Un multiprocesseur à mémoire partagée cesse de grandir vers quelques centaines de cœurs : maintenir un espace d’adressage cohérent au-delà devient trop difficile et trop cher. Les plus gros ordinateurs du monde prennent l’autre chemin. Ce sont des milliers d’ordinateurs ordinaires, chacun avec sa mémoire et son système d’exploitation, qui coopèrent en s’envoyant des messages sur un réseau rapide. Une recherche web, une prévision météo ou l’entraînement d’un grand modèle d’IA tournent tous sur de telles machines.
Ce chapitre porte sur ces multiordinateurs : comment leurs programmes communiquent, ce que coûte un message, comment le réseau est câblé, et à quoi ressemblent les plus grands systèmes actuels.
Le passage de messages
Dans un multiordinateur, aucun processeur ne peut charger une donnée dans la mémoire d’un autre nœud. Si le nœud 3 a besoin d’une donnée détenue par le nœud 7, le nœud 7 doit l’envoyer et le nœud 3 doit la recevoir. La communication est explicite, écrite dans le programme, au lieu d’être cachée derrière des chargements et des rangements.
L’interface de référence du calcul scientifique est MPI (Message Passing Interface), publiée pour la première fois en 1994 et implémentée sur tous les supercalculateurs depuis. Chaque processus d’un calcul a un rang (rank), et les rangs échangent des 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;
}
Le même programme tourne sur chaque nœud, et chaque copie décide quoi faire d’après son rang. Au-delà des messages point à point, MPI fournit des opérations collectives qui impliquent tous les rangs : diffuser une valeur de l’un vers tous, rassembler des résultats, ou faire une réduction globale (all-reduce), où chaque rang apporte un nombre et chaque rang reçoit la somme. La réduction globale est aussi au cœur de l’entraînement des réseaux de neurones sur de nombreux GPU : après chaque étape, les gradients calculés sur chaque GPU sont additionnés sur l’ensemble.
Le passage de messages est plus difficile à programmer que la mémoire partagée, car le programmeur doit décider où vit chaque donnée et quand elle se déplace. Il a un gros avantage : aucun trafic caché. Rien de comparable au faux partage ne peut arriver par accident, et le coût de la communication se voit dans le code.
Ce que coûte un message
Envoyer un message coûte à peu près une latence fixe plus la taille divisée par le débit. La latence, c’est ce que coûte l’envoi d’un seul octet : le logiciel des deux côtés, l’interface réseau, les fils et les commutateurs entre les deux. Voici l’aller-retour d’un tout petit message à différentes distances, mesuré depuis le M2 Ultra :
| Trajet | Aller-retour pour un petit message |
|---|---|
| une ligne de cache entre deux cœurs (mémoire partagée) | 66–400 ns |
| TCP en boucle locale, même machine, macOS | 19–20 µs |
| TCP en boucle locale, VM Linux sur le même Mac | 36–37 µs |
| poignée de main TCP avec le serveur le plus proche d’example.com | 6,4–6,9 ms |
Les chiffres en boucle locale viennent d’un programme C qui fait rebondir 1 octet 20 000 fois entre deux threads à travers une socket TCP sur 127.0.0.1, avec l’algorithme de Nagle désactivé. Aucun fil n’intervient, et c’est pourtant 200 fois plus lent que de déplacer une ligne de cache. Le temps passe dans le noyau : deux appels système par sens, la pile TCP/IP, et le réveil du thread récepteur endormi. La dernière ligne est la durée de la poignée de main TCP pour des requêtes HTTPS vers example.com, telle que la rapporte curl (temps de connexion moins résolution du nom) : un vrai aller-retour sur internet, encore environ 300 fois plus long.
Le débit est l’autre moitié. Le même programme en boucle locale, avec des messages de 1 Mio, a pris 211 à 224 µs par aller-retour, soit environ 9,4 à 9,9 Go/s dans chaque sens. Avec une latence d’environ 10 µs dans un sens et 10 Go/s, un message de 1 Mio est dominé par sa taille, et un message de 8 octets entièrement par la latence. La taille pour laquelle les deux comptent autant est latence × débit, environ 100 Ko ici. Essayez :
À essayer : Appuyez sur Ligne suivante pour exécuter la ligne surlignée, ou sur Lecture pour regarder ; les cases à droite sont les variables, et celles qui viennent de changer s’allument. Modifiez le code pour essayer vos propres changements.
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}
À 10 µs, un message de 4 Ko obtient 3 % du débit maximal et un message de 256 Ko 72 %. Passez latency_ns à 1000, plus proche de ce qu’offre le réseau d’un supercalculateur, et les messages de 32 Ko atteignent déjà 76 %. Les petits messages ne sont que latence : c’est pourquoi les programmes parallèles cherchent à envoyer moins de messages, mais plus gros, et pourquoi les réseaux de calcul haute performance se battent pour chaque microseconde.
Ils y parviennent en écartant le noyau. Des réseaux comme InfiniBand et Slingshot d’HPE utilisent le RDMA (remote direct memory access, accès direct à la mémoire distante) : un processus enregistre une fois un tampon auprès de la carte réseau, après quoi la carte le lit et l’écrit directement, sans appel système ni copie, et peut écrire tout droit dans un tampon du nœud distant. La latence des petits messages entre deux nœuds tombe alors à l’ordre d’une ou deux microsecondes, passées surtout dans les cartes réseau et les commutateurs.
Les réseaux d’interconnexion
Avec des milliers de nœuds, le schéma de câblage (la topologie) compte autant que la vitesse de chaque lien. Deux mesures la décrivent :
- le diamètre : le plus grand nombre de sauts entre deux nœuds, qui borne la latence dans le pire cas ;
- la bande passante de bissection : le nombre de liens à couper pour séparer la machine en deux moitiés, qui borne la quantité de données pouvant traverser le milieu quand chaque nœud parle à un partenaire éloigné.
Pour 64 nœuds :
| Topologie | Liens | Diamètre | Bissection (liens) |
|---|---|---|---|
| anneau | 64 | 32 | 2 |
| maillage 8 × 8 | 112 | 14 | 8 |
| tore 8 × 8 (maillage refermé sur lui-même) | 128 | 8 | 16 |
| hypercube de dimension 6 | 192 | 6 | 32 |
| connexion complète | 2 016 | 1 | 1 024 |
L’anneau est bon marché et sans espoir : la moitié du trafic doit passer par 2 liens. La connexion complète est idéale et impossible au-delà de quelques nœuds. Les vraies machines se situent entre les deux. Les Blue Gene d’IBM utilisaient des tores en 3D puis en 5D ; le japonais Fugaku utilise un maillage/tore à six dimensions appelé Tofu. La plupart des grappes utilisent un arbre épais (fat tree) fait de commutateurs, où les liens s’épaississent (davantage de liens en parallèle) vers la racine, pour que la bande passante de bissection ne rétrécisse pas en haut. Le Slingshot d’HPE utilise une topologie dragonfly : des groupes de commutateurs entièrement connectés à l’intérieur, chaque groupe relié directement à tous les autres, si bien que deux nœuds quelconques ne sont qu’à quelques sauts de commutateur.
La manière dont un message traverse le réseau compte aussi. La commutation stocker-et-retransmettre (store-and-forward) reçoit un paquet entier à chaque commutateur avant de le renvoyer : la latence croît avec le nombre de sauts multiplié par la taille du paquet. Les commutations cut-through et wormhole commencent à retransmettre un paquet dès que son en-tête est arrivé et que la route est décidée : le paquet s’écoule à travers plusieurs commutateurs à la fois, et le surcoût par saut est faible. Tous les réseaux haute performance utilisent la seconde approche.
Des MPP aux grappes
Les premiers grands multiordinateurs étaient des MPP (massively parallel processors, processeurs massivement parallèles) : des milliers de processeurs avec un réseau sur mesure, un conditionnement sur mesure et souvent des puces sur mesure, chez des constructeurs comme Cray, Intel, Thinking Machines ou IBM. Ils étaient chers et liés à un seul fournisseur.
En 1994, un projet de la NASA appelé Beowulf a montré une voie moins chère : une grappe (cluster) de PC ordinaires sous Linux, reliés par de l’Ethernet ordinaire et programmés par passage de messages. Le rapport performance/prix était imbattable, et les grappes de serveurs du commerce ont conquis l’essentiel du calcul haute performance.
La frontière s’est estompée depuis. Les systèmes de tête actuels sont faits de nœuds qui utilisent des processeurs de serveur et des GPU standard et tournent sous Linux, comme le reste de la liste (voir le chapitre sur Unix et Windows), mais ils sont assemblés par quelques constructeurs, avec refroidissement liquide et réseau spécialisé. Ce sont des grappes par l’architecture et des MPP par l’ingénierie.
Les ordinateurs à l’échelle de l’entrepôt
Les plus grandes grappes de toutes ne font pas de simulation scientifique. Une entreprise comme Google, Amazon, Microsoft ou Meta exploite des centres de données de dizaines de milliers de serveurs chacun, et traite un bâtiment entier comme un seul ordinateur qui fait tourner quelques services géants. Des ingénieurs de Google les ont appelés ordinateurs à l’échelle de l’entrepôt (warehouse-scale computers) dans un ouvrage de 2009 consacré au sujet.
Leur conception diffère de celle des supercalculateurs sur trois points :
- La panne est normale. Avec autant de disques, de barrettes mémoire et d’alimentations, quelque chose casse tous les jours. Le logiciel est conçu pour le supporter : les données sont répliquées sur plusieurs machines, et des cadres comme MapReduce (publié par Google en 2004) relancent ailleurs le travail d’une machine tombée en panne. Un calcul sur supercalculateur, à l’inverse, s’arrête en général quand un nœud meurt et repart de son dernier point de reprise (checkpoint).
- Le débit avant la latence. Un moteur de recherche traite des millions de requêtes indépendantes, dont chacune n’a besoin que d’une petite partie de la machine. Ce parallélisme-là est facile ; le difficile, c’est la traîne (tail latency) : quand une requête se déploie sur des centaines de serveurs, c’est le plus lent qui fixe le temps de réponse.
- Énergie et coût. L’alimentation et le refroidissement du bâtiment comptent autant que les serveurs. La mesure de référence, le PUE (power usage effectiveness), divise la puissance totale du site par celle qui atteint les ordinateurs ; les meilleurs centres de données s’approchent de 1, c’est-à-dire que presque toute l’énergie sert au calcul.
L’informatique en nuage (cloud) vend des tranches de ces machines. Elle a largement remplacé l’idée antérieure de grille de calcul (grid computing), qui cherchait à fédérer les grappes de nombreuses institutions en une seule ressource partagée.
Les supercalculateurs et le TOP500
Depuis juin 1993, la liste TOP500 classe deux fois par an, en juin et en novembre, les supercalculateurs les plus rapides du monde. Le classement repose sur un seul test, HPL (High-Performance LINPACK), qui résout un énorme système dense d’équations linéaires en virgule flottante 64 bits. Chaque entrée indique Rmax, la vitesse atteinte sur HPL, et Rpeak, le maximum théorique du matériel.
En juin 2022, Frontier, au laboratoire national d’Oak Ridge aux États-Unis, est devenu le premier système à dépasser l’exaflops (10¹⁸ opérations en virgule flottante par seconde) sur HPL, avec un Rmax de 1,102 exaflops. Il est fait de nœuds HPE Cray EX, chacun combinant un processeur AMD EPYC et des GPU AMD Instinct MI250X, reliés par Slingshot-11. El Capitan, à Lawrence Livermore, a pris la première place en novembre 2024.
Les cinq premiers de la liste de juin 2026, la 67ᵉ édition :
| Rang | Système | Site | Rmax (exaflops) | Puissance (MW) | Processeurs |
|---|---|---|---|---|---|
| 1 | LineShine | NSC Shenzhen, Chine | 2,198 | 42,2 | uniquement des CPU LX2 à 304 cœurs |
| 2 | El Capitan | LLNL, États-Unis | 1,809 | 29,7 | AMD EPYC + MI300A |
| 3 | Frontier | Oak Ridge, États-Unis | 1,353 | 24,6 | AMD EPYC + GPU MI250X |
| 4 | Aurora | Argonne, États-Unis | 1,012 | 38,7 | Intel Xeon Max + GPU Intel |
| 5 | JUPITER Booster | Jülich, Allemagne | 1,000 | 15,8 | NVIDIA GH200 |
Quelques lectures de ce tableau :
- Cinq systèmes exaflopiques. Frontier a été seul pendant deux ans ; ils sont maintenant cinq, et le score de Frontier lui-même est passé de 1,102 à 1,353 à mesure que le système a été optimisé.
- Rmax contre Rpeak. Frontier atteint 66 % de son maximum et El Capitan 64 % ; LineShine 80 %. Même le calcul le plus régulier qui soit n’occupe pas toutes les unités, parce qu’il faut déplacer et échanger des données.
- La puissance. Un exaflops coûte des dizaines de mégawatts, la consommation d’une petite ville. El Capitan fait environ 61 milliards d’opérations par joule, Frontier 55, LineShine 52. C’est l’énergie, et non le nombre de transistors, qui limite désormais la taille de ces machines, pour la même raison que les fréquences ont cessé de monter.
- Les GPU. Quatre des cinq tirent l’essentiel de leur vitesse de GPU ou d’accélérateurs de type GPU, pour les raisons exposées dans le chapitre SIMD et GPU : plus de calcul par watt sur un travail régulier et parallèle sur les données. Le nouveau numéro un est l’exception, construit avec des CPU de 304 cœurs chacun, comme Fugaku l’était avec ses processeurs ARM A64FX quand il menait la liste en 2020 et 2021.
HPL est un test flatteur : de l’arithmétique sur matrices denses, avec beaucoup de calcul par octet déplacé. Beaucoup d’applications réelles sont plutôt limitées par le débit mémoire et la communication. Le test complémentaire HPCG, dominé par des opérations creuses et limitées par la mémoire, donne à LineShine 22,00 pétaflops sur la même liste, environ 1 % de son score HPL.
Faire tourner un programme sur 10 millions de cœurs
La loi d’Amdahl est brutale à cette échelle : avec 10 millions de cœurs, même une fraction séquentielle d’un millionième plafonne l’accélération à environ 900 000. Les supercalculateurs restent utiles parce que les problèmes grandissent avec la machine. Un modèle climatique sur un ordinateur plus gros utilise une grille plus fine, pas la même grille plus vite. Ce passage à l’échelle faible (weak scaling) garde à peu près constante la part de travail de chaque nœud.
La structure habituelle est la décomposition de domaine : l’espace simulé est découpé en blocs, un par processus. Chaque processus met à jour son bloc, puis échange avec ses voisins une fine couche de cellules de bord, le halo. Le calcul croît avec le volume d’un bloc et la communication avec sa surface : des blocs plus gros signifient relativement moins de communication.
Les nœuds sont attribués par un ordonnanceur de traitement par lots comme Slurm : un calcul demande un nombre de nœuds et une durée maximale, attend dans une file, puis obtient ses nœuds pour lui seul pendant toute la durée. Comme les processus d’un calcul parallèle s’attendent les uns les autres à chaque échange, ils doivent tous tourner en même temps ; l’ordonnanceur ne partage jamais les nœuds d’un calcul avec un autre en temps partagé.
On a longtemps essayé de faire ressembler un multiordinateur à une mémoire partagée par le logiciel. Les systèmes de mémoire partagée distribuée, à partir des années 1980, utilisaient le matériel de mémoire virtuelle pour aller chercher une page sur un autre nœud lors d’un défaut de page. Cela fonctionne, mais le faux partage à l’échelle de la page et le coût de chaque défaut les rendaient lents. Le compromis actuel est le modèle PGAS (partitioned global address space, espace d’adressage global partitionné), dans des langages comme UPC et Chapel : un seul espace d’adressage dans le programme, mais chaque donnée a un emplacement visible, si bien que le programmeur sait quels accès sont distants.
À retenir
- Un multiordinateur ou une grappe, ce sont de nombreux ordinateurs à mémoires privées qui coopèrent par passage de messages, en général avec MPI : envoi, réception et opérations collectives comme la réduction globale.
- Coût d’un message ≈ latence + taille / débit. Mesuré depuis le M2 Ultra : 66 à 400 ns pour déplacer une ligne de cache entre cœurs, 19 à 20 µs pour un aller-retour TCP en boucle locale, 6,4 à 6,9 ms jusqu’à un serveur web. La boucle locale a transféré environ 9,4 à 9,9 Go/s avec des messages de 1 Mio.
- Les réseaux haute performance contournent le noyau grâce au RDMA pour ramener la latence entre nœuds à environ une microseconde.
- Les topologies arbitrent entre coût, diamètre et bande passante de bissection : anneaux, maillages, tores, arbres épais, dragonfly. La commutation cut-through garde un faible coût par saut.
- Les grappes de serveurs du commerce (Beowulf, 1994) ont remplacé la plupart des MPP sur mesure. Les ordinateurs à l’échelle de l’entrepôt considèrent la panne comme normale et optimisent le débit, la latence de traîne et l’énergie.
- Le TOP500 classe les systèmes sur HPL. Frontier a franchi 1 exaflops en juin 2022 ; la liste de juin 2026 compte cinq systèmes exaflopiques, menés par LineShine, sans GPU, à 2,198 exaflops, et quatre des cinq reposent sur des GPU. Chacun consomme de 16 à 42 MW.