Site de Jean-Michel RICHER

Maître de Conférences en Informatique à l'Université d'Angers

Ce site est en cours de reconstruction certains liens peuvent ne pas fonctionner ou certaines images peuvent ne pas s'afficher.


stacks

5. MPI

5.1. Introduction

MPI (The Message Passing Interface), conçue en 1993-94, est une norme (ou API - Application Programming Interface) définissant une bibliothèque de fonctions, utilisable avec les langages C, C++ et Fortran. Elle permet d'exploiter des ordinateurs distants ou multiprocesseurs par passage de messages (Wikipedia).

La technique de passage de message consiste à transmettre au travers du réseau les données à échanger. Initialement MPI a été conçu pour des systèmes à mémoire distribuée très populaires dans les années 80-90. MPI fut ensuite adapté pour les systèmes à mémoire partagée et fonctionne à présent pour des systèmes hybrides (distribué + partagé).

L'ensemble des fonctions peut être trouvé à cette adresse OpenMPI Doc.

MPI étant une API, il existe plusieurs implantations come MPICH ou OpenMPI (voir ce site pour un aperçu de l'architecture MPI).

5.2. Installation

Sous Ubuntu 24.04, il n'est plus possible d'utiliser la partie C++ de MPI. Deux possibilités s'offrent à vous :

  • installer OpenMPI à partir des fichiers sources d'une ancienne version (4.1 par exemple)
  • utiliser un archive docker qui utilise OpenMPI

5.2.1. Installation OpenMPI depuis le code source

Récupérez openmpi-4.1.8.tar.gz sur le site d'OpenMPI. Préférer la version 4.1 à la dernière version (5.X) car les fichiers sont moins volumineux.

Décompressez, compilez le code source et installez les librairies, il s'agit d'une installation locale dans votre home directory :

> tar -xvzf openmpi-4.1.8.tar.gz
> cd openmpi-4.1.8/

# en local (pour l'utilisateur courant)
> ./configure --prefix=$HOME/openmpi-4.1.8 --disable-mpi-fortran \
	--disable-oshmem CC=gcc CXX=g++ CFLAGS="-O3 -march=native" \
	CXXFLAGS="-O3 -march=native -std=c++20" --enable-mpi-cxx
> make -j$(nproc)
> make install


# ou pour une installation globale (tous les utilisateurs)
> ./configure prefix=/usr/local/openmpi-4.1.8 --disable-mpi-fortran \
	 --disable-oshmem CC=gcc CXX=g++ CFLAGS="-O3 -march=native" \
	 CXXFLAGS="-O3 -march=native -std=c++20" --enable-mpi-cxx
> make -j$(nproc)
> sudo make install

Note: l'utilisation de -std=c++20 n'est pas forcément nécessaire.

Modifiez ensuite votre fichier .bashrc :

# pour une installation locale
export PATH=\$HOME/openmpi-4.1.8/bin:\$PATH
export LD_LIBRARY_PATH=\$HOME/openmpi-4.1.8/lib:\$LD_LIBRARY_PATH

# ou pour une installation globale
export PATH=/usr/local/openmpi-4.1.8/bin:\$PATH
export LD_LIBRARY_PATH=/usr/local/openmpi-4.1.8/lib:\$LD_LIBRARY_PATH

et faire un source du .bashrc :

> source .bashrc

# vérifier que mpi fonctionne :
> mpirun --version
mpirun (Open MPI) 4.1.8

Report bugs to http://www.open-mpi.org/community/help/

5.2.2. Utilisation de Docker

Commencez par installer Docker s'il n'est pas déjà installé :

> sudo apt-get update
# ne pas oublier docker-buildx pour les versions récentes de docker
> sudo apt-get install docker.io docker-buildx
> sudo usermod -aG docker $USER
> sudo systemctl start docker

Créez ensuite un répertoire dans lequel vous placez un fichier Dockerfile :

> mkdir mpi
> cd mpi
> nano Dockerfile

Le fichier Dockerfile possède le contenu suivant. Il crée un utilisateur user et install OpenMPI ainsi que deux éditeurs en mode texte : nano et fte (commande sfte).

# Utiliser Ubuntu 20.04 comme image de base
FROM ubuntu:20.04

# Empêcher les invites interactives lors de l'installation des paquets
ARG DEBIAN_FRONTEND=noninteractive

# Mettre à jour le système et installer les dépendances nécessaires
# Utilisation de fte pour avoir l'éditeur de texte sfte plus sympathique
# que nano
RUN apt-get update && apt-get install -y \
    build-essential \
    openmpi-bin \
    openmpi-common \
    libopenmpi-dev \
    sudo \
    nano fte fte-console fte-terminal \
    && rm -rf /var/lib/apt/lists/*

# Créer un utilisateur non-root 'user'
RUN useradd -m -s /bin/bash user \
    && echo "user ALL=(ALL) NOPASSWD:ALL" >> /etc/sudoers

# Passer à l'utilisateur 'user'
USER user

# Créer un répertoire de travail pour les programmes
WORKDIR /home/user

# Configurer les variables d'environnement MPI
ENV PATH="/usr/lib64/openmpi/bin:\$PATH" \
    LD_LIBRARY_PATH="/usr/lib64/openmpi/lib:\$LD_LIBRARY_PATH"

# L'utilisateur 'user' peut maintenant compiler et exécuter des programmes MPI
# Vous pouvez ajouter vos fichiers ou configurations ici, par exemple :
# COPY . /home/user/

# Commande par défaut (exemple pour démarrer un shell interactif)
CMD ["/bin/bash"]

Compilez ensuite le Dockerfile afin de créer le container MPI :

> sudo docker build -t mpi_cpp .
> sudo docker run -it mpi_cpp

Ou bien si on veut mapper (mettre en correspondance) le répertoire /home/richer/exchange de l'utilisateur richer sur la machine hôte vers le répertoire /home/user/exchange du container :

> docker run -it \ 
 --mount type=bind,source=/home/richer/exchange,target=/home/user/exchange mpi_cpp

5.3. Utiliser MPI

Il est nécessaire de modifier substantiellement son programme si on désire le paralléliser avec MPI car on exécute le même programme sur plusieurs machines/coeurs différents.

MPI est disponible pour Fortran, C et C++ et également avec Python (pyMPI).

Pour pouvoir disposer de MPI sur sa machine il faut installer les librairies suivantes sous Ubuntu :

> sudo apt-get install libopenmpi-dev openmpi-bin openmpi-doc

Pour C et C++ on inclura le fichier mpi.h

5.3.1. Débogage

Pour débuguer il est recommandé d'installer MPE (Multi-Processing Environment), difficile à trouver, qui doit être compilé lors de l'installation de Python Anaconda. On réalisera alors l'édition de liens avec -llmpe -lmpe

On peut également utiliser la commande suivante (ici dans le cas de deux programmes avec gdb):

> mpirun -n 2 xterm -e gdb ./a.exe

Deux terminaux sont alors ouverts et il faut lancer l'exécution du programme dans le débogueur avec run.

5.4. Communicateur

Les opérations MPI portent sur des communicateurs qui sont en quelque sorte des ports de communication inter processus. Le communicateur par défaut est COMM_WORLD qui comprend tous les processus actifs.

On peut créer de nouveaux communicateurs mais cela se révèle complexe.

Le modèle de communication est un modèle point à point basé sur deux opérations élémentaires send et receive qui sont ensuite déclinés de plusieurs manières différentes.

  • broadcast : envoi une ou plusieurs données d'un membre vers les autres membres du groupe
  • gather : rassembler des données des autres membres vers un membre pour en faire un tableau
  • scatter : découpage d'un tableau et envoi de chaque partie du tableau aux autres membres du groupe
  • reduce : algorithme de réduction
  • scan : algorithme de scan

Le schéma classique d'utilisation de MPI consiste à :

  • initialiser les ressources MPI_Init
  • obtenir le nombre de processus lancés MPI_Comm_size
  • obtenir l'identifiant du processus courant MPI_Comm_rank
  • exécuter le code séquentiel
  • exécuter le code parallèle
  • libérer les ressources MPI_Finalize

Voici deux exemples simples qui affichent un message avec un temps de latence (instruction sleep) pour chaque thread / processeur. Le premier utilise des fonctions C, le second utilise des espaces de noms et des objets

Afficher le code    ens/m2/paracpu/mpi_c_basic.cpp
  1. // ==================================================================
  2. // Author: Jean-Michel Richer
  3. // Email: jean-michel.richer@univ-angers.fr
  4. // Date: Aug 2016
  5. // Purpose: Demonstrate basic functionnalities of MPI with C
  6. // ==================================================================
  7. #include <iostream>
  8. #include <cstdlib>
  9. #include <unistd.h> // for sleep
  10. using namespace std;
  11. #include <mpi.h>
  12.  
  13. // ==============================================================
  14. // C version
  15. // ==============================================================
  16.  
  17. int main(int argc, char ** argv) {
  18.     // maximum number of CPUs
  19.     int max_cpus;
  20.     // cpu identifier (called rank)
  21.     int cpu_rank;
  22.     // C-string to store the name of the host
  23.     char cpu_name[MPI_MAX_PROCESSOR_NAME];
  24.     // length of the C-string
  25.     int length;
  26.  
  27.  
  28.     // initialization also using command line parameters
  29.     MPI_Init(&argc, &argv);
  30.    
  31.     // get number of programs running
  32.     MPI_Comm_size(MPI_COMM_WORLD, &max_cpus);
  33.  
  34.     // get program identifier
  35.     MPI_Comm_rank(MPI_COMM_WORLD, &cpu_rank);
  36.  
  37.     // get cpu name
  38.     MPI_Get_processor_name(cpu_name, &length);
  39.  
  40.     sleep(cpu_rank);
  41.    
  42.     cerr << "running on " << cpu_name << " with id=" << cpu_rank << "/";
  43.     cerr << max_cpus << endl;
  44.    
  45.     // free resources, don't forget !
  46.     MPI_Finalize();
  47.  
  48.     exit(EXIT_SUCCESS);
  49. }
  50.  
Afficher le code    ens/m2/paracpu/mpi_cpp_basic.cpp
  1. // ==================================================================
  2. // Author: Jean-Michel Richer
  3. // Email: jean-michel.richer@univ-angers.fr
  4. // Date: Aug 2016
  5. // Purpose: Demonstrate basic functionnalities of MPI with C++
  6. // ==================================================================
  7. #include <iostream>
  8. #include <cstdlib>
  9. #include <unistd.h> // for sleep
  10. using namespace std;
  11. #include <mpi.h>
  12.  
  13. // ==================================================================
  14. // C++ version
  15. // ==================================================================
  16.  
  17. int main(int argc, char ** argv) {
  18.     // maximum number of CPUs
  19.     int max_cpus;
  20.     // cpu identifier (called rank)
  21.     int cpu_rank;
  22.     // C-string to store name of the host
  23.     char cpu_name[MPI::MAX_PROCESSOR_NAME];
  24.     // length of the C-string
  25.     int length;
  26.  
  27.     // initialization also using command line parameters
  28.     MPI::Init(argc, argv);
  29.    
  30.     // get number of programs running
  31.     max_cpus = MPI::COMM_WORLD.Get_size();
  32.    
  33.     // get program identifier
  34.     cpu_rank = MPI::COMM_WORLD.Get_rank();
  35.    
  36.     // get cpu name
  37.     memset(cpu_name, 0, MPI::MAX_PROCESSOR_NAME);
  38.     MPI::Get_processor_name(cpu_name, length);
  39.  
  40.     // sleep
  41.     sleep(cpu_rank);
  42.    
  43.     cerr << "running on " << cpu_name << " with id=" << cpu_rank << "/";
  44.     cerr << max_cpus << endl;
  45.    
  46.     // free resources, don't forget !
  47.     MPI::Finalize();
  48.  
  49.     exit(EXIT_SUCCESS);
  50. }
  51.  
  52.  

5.4.1. Compilation

Afin de compiler un programme MPI, on utilise le compilateur mpicc pour le C, mpic++ pour le C++ :

# pour le code source C
> mpicc -o exe  src.c -O3 ... 
# ou alors pour les sources en C++ :
> mpic++ -o exe  src.cpp -O3 ... 
# ou alors si cela ne fonctionne pas :
> mpicxx -o exe  src.cpp -O3 ... 

5.4.2. Execution

Pour exécuter un programme compilé avec MPI, il faut utiliser mpirun (ou mpiexec) en spécifiant le nombre de processus (=machines, processeurs, coeurs) grâce à l'option -n (number) ou -np (number of processes) suivant les systèmes :

> mpirun -n 4 ./c3_mpi_ex_1.exe

Dans le cas où on voudrait lancer plus d'instances de programme que de threads disponibles, mpirun refusera d'exécuter le programme :

> mpirun -n 20 ./ezmpi_broadcast.exe
--------------------------------------------------------------------------
There are not enough slots available in the system to satisfy the 20
slots that were requested by the application:
...
Alternatively, you can use the --oversubscribe option to ignore the
number of available slots when deciding the number of processes to
launch
--------------------------------------------------------------------------

Il faudra alors utiliser l'option mpirun --oversubscribe -n 20 ....

Voici le résultat à l'affichage du programme précédent si on l'exécute sur une seule machine Intel Core i3-2375M CPU @ 1.50GHz (Dual Core + HyperThreading, soit 4 threads) :

running on inspiron with id=0/4
running on inspiron with id=1/4
running on inspiron with id=2/4
running on inspiron with id=3/4

Voici un autre exemple qui lance 16 processus sur 4 noeuds (h1 à h4) en réservant 4 processus (coeurs) sur chacune des machines :

> mpiexec -hosts h1:4,h2:4,h3:4,h4:4 -n 16 ./test

Référez vous à la section 5.10 pour voir réellement comment faire.

Sur un cluster on utilisera par exemple (-pe = parallel environment):

> qsub -pe mpi 16 test_mpi.sh

5.5. Synchronisation

Pour obtenir une section critique pour l'affichage on utilise la fonction MPI_Barrier ou MPI::COMM_WORLD.Barrier() en C++ qui bloque l'appelant jusqu'à ce que tous les autres programmes aient appelé cette fonction :

Afficher le code    ens/m2/paracpu/mpi_cpp_synchro.cpp
  1. // ==================================================================
  2. // Author: Jean-Michel Richer
  3. // Email: jean-michel.richer@univ-angers.fr
  4. // Date: Aug 2016
  5. // Purpose: Demonstrate basic functionalities of MPI with C++
  6. // with synchronization for display
  7. // ==================================================================
  8. #include <iostream>
  9. #include <cstdlib>
  10. #include <unistd.h> // for sleep
  11. using namespace std;
  12. #include <mpi.h>
  13.  
  14. // ==================================================================
  15. // C++ version
  16. // ==================================================================
  17.  
  18. int main(int argc, char ** argv) {
  19.     // maximum number of CPUs
  20.     int max_cpus;
  21.     // cpu identifier (called rank)
  22.     int cpu_rank;
  23.     // C-string to store the name of the host
  24.     char cpu_name[MPI_MAX_PROCESSOR_NAME];
  25.     // length of the C-string
  26.     int length;
  27.  
  28.  
  29.     // initialization also using command line parameters
  30.     MPI_Init(&argc, &argv);
  31.    
  32.     // get number of programs running
  33.     max_cpus = MPI::COMM_WORLD.Get_size();
  34.    
  35.     // get program identifier
  36.     cpu_rank = MPI::COMM_WORLD.Get_rank();
  37.    
  38.     // get cpu name
  39.     memset(cpu_name, 0, MPI::MAX_PROCESSOR_NAME);
  40.     MPI::Get_processor_name(cpu_name, length);
  41.  
  42.     cout << "Hello, from " << cpu_rank << endl;
  43.    
  44.     sleep( cpu_rank + 1 );
  45.    
  46.     MPI::COMM_WORLD.Barrier();    
  47.     cout << "Bye, from " << cpu_rank << endl;
  48.    
  49.     // free resources, don't forget !
  50.     MPI::Finalize();
  51.  
  52.     exit(EXIT_SUCCESS);
  53. }
  54.  
  55.  

Dans l'exemple qui suit, les 6 processus commencent par travailler (travail remplacé ici par un sleep) en parallèle. Le premier (identifiant 0) suspend son exécution pendant 1 seconde, le second pendant 2 secondes, etc. Au final, après 6 secondes, l'ensemble des processus exécutent le code qui suit l'appel à MPI::COMM_WORLD.Barrier();.

> mpirun -n 6 ./mpi_cpp_syncrho.exe
Hello, from 3
Hello, from 4
Hello, from 5
Hello, from 2
Hello, from 0
Hello, from 1
.... wait 6 seconds here until you finally get ....
Bye, from 0
Bye, from 5
Bye, from 3
Bye, from 4
Bye, from 1
Bye, from 2

5.6. Echange de donneés Send, Recv

Pour transmettre des données entre les processeurs, on utilise deux fonctions MPI_Send et MPI_Recv, ou éventuellement MPI_Sendrecv :

// envoi d'un processus vers un autre processus
int MPI_Send(void* data,
    int count,
    MPI_Datatype datatype,
    int destination,
    int tag,
    MPI_Comm communicator)

// réception d'un processus vers un autre processus
int MPI_Recv(void* data,
    int count,
    MPI_Datatype datatype,
    int source,
    int tag,
    MPI_Comm communicator,
    MPI_Status* status)

// envoi suivi d'une réception    
int MPI_Sendrecv(const void *sendbuf, 
    int sendcount, MPI_Datatype sendtype,
    int dest, int sendtag,
    void *recvbuf, int recvcount, MPI_Datatype recvtype,
    int source, int recvtag,
    MPI_Comm comm, MPI_Status *status)    
data
pointeur sur la donnée à envoyer ou recevoir
count
nombre d'occurrences
datatype
type de donnée (MPI_CHAR, MPI_INT, MPI_FLOAT, MPI::INT, MPI::FLOAT, ...)
source/destination
identifiant du processus qui reçoit ou envoie les données
communicator
canal de comunication, ex: COMM_WORLD
tag
identifiant de message compris entre 0 et 32767
status
le status lors de la réception des données

Pour la partie C++ on utilisera :

MPI::COMM_WORLD.Send(const void* buf, 
	int count, 
	MPI::Datatype& datatype,
	int dest, 
	int tag) const
	
void MPI::COMM_WORLD.Recv(void* buf, 
	int count, 
	MPI::Datatype& datatype,
	int source, 
	int tag, 
	MPI::Status* status) const
	
void MPI::COMM_WORLD::Sendrecv(const void* sendbuf, 
	int count,
	MPI::Datatype& datatype, 
	int dest, int sendtag,
	void* recvbuf, int recvcount,
	MPI::Datatype recvtype, int source, int recvtag,
	MPI::Status* status) const

Voici quelques exemples :

Dans le premier exemple, on utilise deux processeurs. Le maître (identifiant 0) et l'esclave (identifiant 1) :

 Maître   Esclave 
 Création d'un tableau et initialisation    
 Envoi de la taille →    
    Réception de la taille, création du tableau 
 Envoi du tableau →    
    Réception des données dans le tableau 
    Calcul de la somme des éléments 
    ← Envoi de la somme 
 Réception de la somme    
 Affichage de la somme    
Fonctionnement Exemple Send et Receive
Afficher le code    ens/m2/paracpu/mpi_send_and_recv.cpp
  1. // ==================================================================
  2. // Author: Jean-Michel Richer
  3. // Email: jean-michel.richer@univ-angers.fr
  4. // Date: Aug 2016
  5. // Purpose: Demonstrate basic send functionality, send array of float
  6. // from master (rank=0) to slave (rank=1)
  7. // ==================================================================
  8. #include <unistd.h>
  9. #include <iostream>
  10. #include <cstdlib>
  11. #include <unistd.h> // for sleep
  12. using namespace std;
  13. #include <mpi.h>
  14.  
  15. // ==============================================================
  16. // C++ version
  17. // ==============================================================
  18.  
  19. int main(int argc, char ** argv) {
  20.     // maximum number of CPUs
  21.     int max_cpus;
  22.     // cpu identifier (called rank)
  23.     int cpu_rank;
  24.     // C-string to store name of the host
  25.     char cpu_name[MPI::MAX_PROCESSOR_NAME];
  26.     // length of the C-string
  27.     int length;
  28.  
  29.     // initialization also using command line parameters
  30.     MPI::Init(argc, argv);
  31.  
  32.     // get number of programs running
  33.     max_cpus = MPI::COMM_WORLD.Get_size();
  34.    
  35.     // get program identifier
  36.     cpu_rank = MPI::COMM_WORLD.Get_rank();
  37.    
  38.     // get cpu name
  39.     memset(cpu_name, 0, MPI::MAX_PROCESSOR_NAME);
  40.     MPI::Get_processor_name(cpu_name, length);
  41.  
  42.    
  43.     cerr << cpu_rank << "/" << max_cpus << " on machine " << cpu_name << endl;
  44.    
  45.     // number of elements of the array 'data' that will be allocated
  46.     int data_size;
  47.     // array of float created by the master and sent to the slave(s)
  48.     float *data;
  49.     // identifier (rank) of remove processor (master or slave)
  50.     int remote_cpu;
  51.    
  52.     // Status for communication
  53.     MPI::Status status;
  54.  
  55.     // MASTER
  56.     if (cpu_rank == 0) {
  57.         // MASTER
  58.         // processor of id 0 (master) fills and sends the array
  59.         data_size = 100 + rand() % 100;
  60.         data = new float [data_size];
  61.         for (int i=0; i<data_size; ++i) data[i] = i+1;
  62.        
  63.         int remote_cpu = 1;
  64.        
  65.         // 1- send length
  66.         cerr << cpu_rank << " [send]  data_size=" << data_size << endl;
  67.         MPI::COMM_WORLD.Send(&data_size, 1, MPI::INT, remote_cpu, 0);
  68.         // 2- send data
  69.         MPI::COMM_WORLD.Send(&data[0], data_size, MPI::FLOAT, remote_cpu, 0);      
  70.         // wait for other processor to compute and send back sum
  71.         float result = 0;
  72.  
  73.         MPI::COMM_WORLD.Recv(&result, 1, MPI::FLOAT, remote_cpu, MPI::ANY_TAG, status);
  74.         cerr << cpu_rank << " [recv] result=" << result << endl;
  75.         cerr << "result is " << result << " for length=" << data_size;
  76.         cerr << ", expected=" << (data_size * (data_size+1))/ 2  << endl;
  77.         delete [] data;
  78.        
  79.     } else if (cpu_rank == 1) {
  80.         // SLAVE
  81.         // processor of id 1 will receive the array and compute the sum
  82.         // 1- we receive the array length and allocate space
  83.         remote_cpu = 0; // master
  84.         MPI::COMM_WORLD.Recv(&data_size, 1, MPI::INT, remote_cpu, MPI::ANY_TAG, status);
  85.         cerr << cpu_rank << " [recv] data_size=" << data_size << endl;
  86.         data = new float[data_size];
  87.         // 2- we receive the data
  88.         MPI::COMM_WORLD.Recv(&data[0], data_size, MPI::FLOAT, remote_cpu, MPI::ANY_TAG, status);
  89.         // 3- compute sum
  90.         float sum = 0;
  91.         for (int i=0; i<data_size; ++i) {
  92.             sum += data[i];
  93.         }
  94.         // 4- send back result
  95.         cerr << cpu_rank << " [send]  sum=" << sum << endl;
  96.         MPI::COMM_WORLD.Send(&sum, 1, MPI::FLOAT, remote_cpu, 0);
  97.         delete [] data;
  98.        
  99.     } else {
  100.         // other processors (if any) won't do anything
  101.     }
  102.    
  103.     MPI::Finalize();
  104.  
  105.     exit(EXIT_SUCCESS);
  106. }
  107.  

Le deuxième exemple utilise sendrecv pour échanger une donné entre maître et esclave.

Attention la fonction Sendrecv permet uniquement d'échanger deux données de même type et de même nombre d'occurrence.

On ne peut pas, par exemple, envoyer un tableau de données et attendre la somme en retour.

Afficher le code    ens/m2/paracpu/mpi_sendrecv.cpp
  1. #include <unistd.h>
  2. #include <iostream>
  3. #include <cstdlib>
  4. #include <unistd.h> // for sleep
  5. #include <sstream>
  6. #include <numeric>
  7. using namespace std;
  8. #include <mpi.h>
  9.  
  10.  
  11. // ==============================================================
  12. // version C++
  13. // ==============================================================
  14.  
  15. int main(int argc, char ** argv) {
  16.     // maximum number of CPUs
  17.     int max_cpus;
  18.     // cpu identifier (called rank)
  19.     int cpu_rank;
  20.     // C-string to store name of the host
  21.     char cpu_name[MPI::MAX_PROCESSOR_NAME];
  22.     // length of the C-string
  23.     int length;
  24.  
  25.     // initialization also using command line parameters
  26.     MPI::Init(argc, argv);
  27.  
  28.     // get number of programs running
  29.     max_cpus = MPI::COMM_WORLD.Get_size();
  30.    
  31.     // get program identifier
  32.     cpu_rank = MPI::COMM_WORLD.Get_rank();
  33.    
  34.     // get cpu name
  35.     memset(cpu_name, 0, MPI::MAX_PROCESSOR_NAME);
  36.     MPI::Get_processor_name(cpu_name, length);
  37.  
  38.     cerr << cpu_rank << "/" << max_cpus << " on machine " << cpu_name << endl;
  39.    
  40.     // data to exchange
  41.     int send_data;
  42.     int recv_data;
  43.    
  44.     // identifier (rank) of remove processor (master or slave)
  45.     int remote_cpu, source_cpu;
  46.    
  47.     // Status for communication
  48.     MPI::Status status;
  49.     int TAG_0 = 0;
  50.  
  51.     if (cpu_rank == 0) {
  52.        
  53.         remote_cpu = 1;
  54.         send_data = 11111;
  55.  
  56.     } else if (cpu_rank == 1) {
  57.        
  58.         remote_cpu = 0;
  59.         send_data = 22222;
  60.        
  61.     }
  62.  
  63.     MPI::COMM_WORLD.Sendrecv(
  64.         &send_data,   // Send buffer
  65.         1,            // Number of elements to send
  66.         MPI::INT,     // Datatype of send elements
  67.         remote_cpu,   // Destination rank
  68.         0,            // Send tag
  69.         &recv_data,   // Receive buffer
  70.         1,            // Number of elements to receive
  71.         MPI::INT,     // Datatype of receive elements
  72.         remote_cpu,   // Source rank
  73.         0,            // Receive tag
  74.         status        // Status object
  75.     );
  76.    
  77.    
  78.     if (cpu_rank == 0) {
  79.        
  80.         cerr << "Master: recv_data=" << recv_data << endl;
  81.        
  82.     } else if (cpu_rank == 1) {
  83.        
  84.         cerr << "Slave: recv_data=" << recv_data << endl;
  85.  
  86.     }
  87.     MPI::Finalize();
  88.  
  89.     exit(EXIT_SUCCESS);
  90. }
  91.  

Afin de simplifier l'utilisation de MPI on peut créer une interface (wrapper) : j'ai créé, pour ma part, EZMPI.

EZMPI est constitué d'une seule classe Process qui représente un processus et qui permet de récupérer lors de l'initialisation l'identifiant du processeur, le nombre de processus lancés, ... On dispose également de méthodes qui permettent d'envoyer et recevoir un élément (quel que soit le type) ou un tableau d'élements.

On dispose en outre d'un système de log qui permet de consulter le déroulement des échanges entre processus (réception et envoi de données).

Le fichier principal à inclure est le suivant (il comprend interface et implantation) :

Afficher le code    ens/m2/paracpu/ezmpi.hpp
  1. /*
  2. This file is part of the EZLib Library
  3. Copyright (C) 2026 Jean-Michel RICHER
  4.  
  5. This program is free software: you can redistribute it and/or modify
  6. it under the terms of the GNU General Public License as published by
  7. the Free Software Foundation, either version 3 of the License, or
  8. (at your option) any later version.
  9.  
  10. This program is distributed in the hope that it will be useful,
  11. but WITHOUT ANY WARRANTY; without even the implied warranty of
  12. MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the
  13. GNU General Public License for more details.
  14.  
  15. You should have received a copy of the GNU General Public License
  16. along with this program.  If not, see <https://www.gnu.org/licenses/>.
  17. */
  18.  
  19. // ==================================================================
  20. // Author: Jean-Michel Richer
  21. // Email: jean-michel.richer@univ-angers.fr
  22. // Date: Aug 2020
  23. // Last modified: September 2026
  24. // Purpose: interface and class to facilitate use of MPI
  25. // ==================================================================
  26. #pragma once
  27.  
  28. #include <time.h>
  29. #include <unistd.h>
  30.  
  31. #include <cassert>
  32. #include <cstdint>
  33. #include <sstream>
  34. #include <stdexcept>
  35. #include <string>
  36. #include <typeinfo>
  37. #include <vector>
  38. using namespace std;
  39. #include <mpi.h>
  40.  
  41. /**
  42.  * @brief EZ MPI is a wrapper for MPI C++. It simplifies the use
  43.  * of MPI send, receive, gather, scatter functions.
  44.  * For the gather, scatter and reduce functions the processor
  45.  * of rank 0 is considered as the "master" that collects or
  46.  * sends data.
  47.  */
  48.  
  49. namespace ez {
  50.  
  51. namespace mpi {
  52.  
  53. enum { MASTER = 0, SLAVE_1, SLAVE_2, SLAVE_3, SLAVE_4 };
  54.  
  55. /**
  56.  * @class Process
  57.  * @brief This is the main class that defines a process.
  58.  * This class gathers information about the number of other processes
  59.  * working with this process and the rank of the process also called
  60.  * id here.
  61.  * It enables to perform operations like send and receive but also more
  62.  * complex ones as gather, scatter, broadcast or reduce which are supposed
  63.  * to be initiated by the master process which as the default rank of 0.
  64.  */
  65. class Process {
  66.    protected:
  67.     // number of process working with this one
  68.     int m_nbr_processes;
  69.  
  70.     // identifier of current process, called rank for MPI
  71.     int m_id;
  72.  
  73.     // identifier of remote process for to communicate with (send, receive, ...)
  74.     int m_remote;
  75.  
  76.     // status of last operation
  77.     MPI::Status m_status;
  78.  
  79.     // message tag if needed (default is 0)
  80.     int m_message_tag;
  81.  
  82.     // name of process
  83.     string m_name;
  84.  
  85.     // Linux process identifier (PID)
  86.     int m_pid;
  87.  
  88.     // verbose mode
  89.     bool m_verbose_flag;
  90.  
  91.     // verbose mode
  92.     bool m_log_flag;
  93.  
  94.     // main output stream for current processor
  95.     ostringstream log_stream;
  96.  
  97.     // temporary output stream
  98.     ostringstream tmp_log;
  99.  
  100.     // MPI::Finalize already called ?
  101.     bool finalize_already_called;
  102.  
  103.    private:
  104.     /**
  105.      * Find processor name
  106.      */
  107.     void find_cpu_name() {
  108.         char name[MPI::MAX_PROCESSOR_NAME];
  109.         int length;
  110.  
  111.         memset(name, 0, MPI::MAX_PROCESSOR_NAME);
  112.         MPI::Get_processor_name(name, length);
  113.         m_name = name;
  114.     }
  115.  
  116.     /**
  117.      * record output of oss into general output
  118.      */
  119.     // void append();
  120.  
  121.     /**
  122.      * initialize nbr_process, cpu_rank, processus id
  123.      */
  124.     void init() {
  125.         m_nbr_processes = MPI::COMM_WORLD.Get_size();
  126.         m_id = MPI::COMM_WORLD.Get_rank();
  127.         find_cpu_name();
  128.         m_pid = getpid();
  129.         if (m_verbose_flag) {
  130.             tmp_log << "pid=" << m_pid << ", id=" << m_id << endl;
  131.             flush();
  132.         }
  133.     }
  134.  
  135.     const std::string current_date_and_time() {
  136.         time_t now = time(0);
  137.         struct tm tstruct;
  138.         char buf[80];
  139.         tstruct = *localtime(&now);
  140.         // strftime(buf, sizeof(buf), "%Y-%m-%d.%X [%s]", &tstruct);
  141.         strftime(buf, sizeof(buf), "%X", &tstruct);
  142.  
  143.         return buf;
  144.     }
  145.  
  146.     void flush() {
  147.         string str = current_date_and_time();
  148.         if (m_log_flag) {
  149.             log_stream << str << " cpu " << m_id << "/" << m_nbr_processes
  150.                        << ": " << tmp_log.str();
  151.         }
  152.         if (m_verbose_flag) {
  153.             cerr << str << " cpu " << m_id << "/" << m_nbr_processes << ": "
  154.                  << tmp_log.str();
  155.         }
  156.         tmp_log.str("");
  157.     }
  158.  
  159.     void print(char v) { tmp_log << v; }
  160.     void print(int v) { tmp_log << v; }
  161.     void print(string v) { tmp_log << v; }
  162.     void print(float v) { tmp_log << v; }
  163.     void print(double v) { tmp_log << v; }
  164.  
  165.    public:
  166.     /**
  167.      * Default constructor
  168.      * @param argc command line number of arguments
  169.      * @param argv command line arguments
  170.      * @param verbose allow verbose mode
  171.      * @param log allow logging mode
  172.      */
  173.     Process(int argc, char* argv[], bool verbose = true, bool log = true) {
  174.         m_remote = 0;
  175.         m_message_tag = 0;
  176.         m_verbose_flag = verbose;
  177.         m_log_flag = log;
  178.         finalize_already_called = false;
  179.         MPI::Init(argc, argv);
  180.         init();
  181.     }
  182.  
  183.     /**
  184.      * Free resources of MPI::Init
  185.      */
  186.     void finalize() {
  187.         if (!finalize_already_called) {
  188.             finalize_already_called = true;
  189.             MPI::COMM_WORLD.Barrier();
  190.             MPI::Finalize();
  191.         }
  192.     }
  193.  
  194.     ~Process() {
  195.         if (!finalize_already_called) {
  196.             finalize_already_called = true;
  197.             MPI::COMM_WORLD.Barrier();
  198.             MPI::Finalize();
  199.         }
  200.     }
  201.  
  202.     /**
  203.      * set verbose mode
  204.      * @param mode true to enable, false to disable
  205.      */
  206.     void verbose(bool mode) { m_verbose_flag = mode; }
  207.  
  208.     /**
  209.      * set log mode
  210.      * @param mode true to enable, false to disable
  211.      */
  212.     void log(bool mode) { m_log_flag = mode; }
  213.  
  214.     /**
  215.      * Get Unix process identifier
  216.      * @return unix process identifier
  217.      */
  218.     int unix_pid() { return m_pid; }
  219.  
  220.     /**
  221.      * Get processor identifier or rank
  222.      * @return rank of process
  223.      */
  224.     int rank() { return m_id; }
  225.  
  226.     /**
  227.      * Get processor identifier or rank
  228.      * @return rank of process
  229.      * @note this method is identical to rank()
  230.      */
  231.     int id() { return m_id; }
  232.  
  233.     /**
  234.      * Get number of processors used
  235.      * @return number of other processes working with this process
  236.      */
  237.     int nbr_processes() { return m_nbr_processes; }
  238.  
  239.     /**
  240.      * Get processor name
  241.      * @return name of processor on which this process is executed
  242.      */
  243.     string name() { return m_name; }
  244.  
  245.     /**
  246.      * set remote processor identifier
  247.      * @param rmt_id remote process identifier
  248.      */
  249.     void remote(int rmt_id) { m_remote = rmt_id; }
  250.  
  251.     /**
  252.      * set message tag
  253.      * @param tag must be an integer between 0 and 32767
  254.      */
  255.     void tag(int tag) { m_message_tag = tag; }
  256.  
  257.     /**
  258.      * Return true if this processor is the processor of rank 0
  259.      * considered as the master.
  260.      */
  261.     bool is_master() { return (m_id == ez::mpi::MASTER); }
  262.  
  263.     bool is_slave() { return (m_id != ez::mpi::MASTER); }
  264.  
  265.     bool is_slave(int n) { return (m_id == n); }
  266.  
  267.     /**
  268.      * synchronize
  269.      */
  270.     void synchronize() { MPI::COMM_WORLD.Barrier(); }
  271.  
  272.     /**
  273.      * Determine type of data T and convert it into MPI::Datatype.
  274.      * This function needs to be extended with other types.
  275.      */
  276.     template <class T>
  277.     MPI::Datatype get_type() {
  278.         if (typeid(T) == typeid(char)) {
  279.             return MPI::CHAR;
  280.         } else if (typeid(T) == typeid(int8_t)) {
  281.             return MPI::CHAR;
  282.         } else if (typeid(T) == typeid(uint8_t)) {
  283.             return MPI::CHAR;
  284.         } else if (typeid(T) == typeid(int)) {
  285.             return MPI::INT;
  286.         } else if (typeid(T) == typeid(float)) {
  287.             return MPI::FLOAT;
  288.         } else if (typeid(T) == typeid(double)) {
  289.             return MPI::DOUBLE;
  290.         }
  291.         std::stringstream oss;
  292.         oss << "!!!!! unknown " << typeid(T).name() << endl;
  293.         throw std::runtime_error(oss.str());
  294.  
  295.         return MPI::INT;
  296.     }
  297.  
  298.     /**
  299.      * Determine type of data T and convert it into a string
  300.      */
  301.     template <class T>
  302.     std::string get_type_name() {
  303.         if (typeid(T) == typeid(char)) {
  304.             return "char";
  305.         } else if (typeid(T) == typeid(int8_t)) {
  306.             return "char";
  307.         } else if (typeid(T) == typeid(uint8_t)) {
  308.             return "char";
  309.         } else if (typeid(T) == typeid(int)) {
  310.             return "int";
  311.         } else if (typeid(T) == typeid(float)) {
  312.             return "float";
  313.         } else if (typeid(T) == typeid(double)) {
  314.             return "double";
  315.         }
  316.         std::stringstream oss;
  317.         oss << "!!!!! unknown " << typeid(T).name() << endl;
  318.         throw std::runtime_error(oss.str());
  319.  
  320.         return "char";
  321.     }
  322.  
  323.     /**
  324.      * Send one instance of data to remote_cpu
  325.      * @param v data to send
  326.      */
  327.     template <class T>
  328.     void send(T& v) {
  329.         MPI::Datatype data_type = get_type<T>();
  330.  
  331.         if (m_log_flag) {
  332.             tmp_log << "- send value=" << v << " of type "
  333.                     << get_type_name<T>();
  334.             tmp_log << " to process=" << m_remote << endl;
  335.             flush();
  336.         }
  337.  
  338.         MPI::COMM_WORLD.Send(&v, 1, data_type, m_remote, m_message_tag);
  339.     }
  340.  
  341.     /**
  342.      * Send an array to remote_cpu
  343.      * @param arr address of the array
  344.      * @param size number of elements to send
  345.      */
  346.     template <class T>
  347.     void send(T* arr, int size, int begin = 0, int end = -1) {
  348.         MPI::Datatype data_type = get_type<T>();
  349.  
  350.         assert((begin >= 0) and (begin < size));
  351.         if (end < 0) end = size;
  352.         assert((end >= 0) and (end <= size));
  353.         size = end - begin;
  354.  
  355.         if (m_log_flag) {
  356.             tmp_log << "- send array of type " << get_type_name<T>();
  357.             tmp_log << " from index " << begin << " to " << end;
  358.             tmp_log << " of size " << size << " to process=" << m_remote
  359.                     << endl;
  360.             flush();
  361.         }
  362.  
  363.         MPI::COMM_WORLD.Send(&arr[begin], size, data_type, m_remote,
  364.                              m_message_tag);
  365.     }
  366.  
  367.     /**
  368.      * Send data from size data of the vector starting at
  369.      * given index
  370.      * @param v vector
  371.      * @param size number of elements to send
  372.      * @param index where to start, default is first
  373.      *        element of the vector
  374.      *
  375.      */
  376.     template <class T>
  377.     void send(vector<T>& v, int begin = 0, int end = -1) {
  378.         int size = static_cast<int>(v.size());
  379.         assert((begin >= 0) and (begin < size));
  380.         if (end < 0) end = size;
  381.         assert((end >= 0) and (end <= size));
  382.         size = end - begin;
  383.  
  384.         MPI::Datatype data_type = get_type<T>();
  385.  
  386.         MPI::COMM_WORLD.Send(&size, 1, MPI::INT, m_remote, m_message_tag);
  387.         T* ptr = v.data();
  388.         MPI::COMM_WORLD.Send(&ptr[begin], size, data_type, m_remote,
  389.                              m_message_tag);
  390.  
  391.         if (m_log_flag) {
  392.             tmp_log << "- send " << size << " elements of vector";
  393.             tmp_log << " of type " << get_type_name<T>();
  394.             tmp_log << " from index " << begin << " to " << end;
  395.             tmp_log << " of size " << size << " to process=" << m_remote
  396.                     << endl;
  397.             flush();
  398.         }
  399.     }
  400.  
  401.     /**
  402.      * Receive one instance of data from remote cpu
  403.      * @param v data to receive
  404.      */
  405.     template <class T>
  406.     void recv(T& v) {
  407.         MPI::Datatype data_type = get_type<T>();
  408.  
  409.         MPI::COMM_WORLD.Recv(
  410.             &v, 1, data_type, m_remote,
  411.             (m_message_tag == 0) ? MPI::ANY_TAG : m_message_tag, m_status);
  412.  
  413.         if (m_log_flag) {
  414.             tmp_log << "- receive value=" << v << " of type "
  415.                     << get_type_name<T>();
  416.             tmp_log << " from process=" << m_remote << endl;
  417.             flush();
  418.         }
  419.     }
  420.  
  421.     /**
  422.      * Receive an array of given size
  423.      * @param arr pointer to address of the array
  424.      * @param size number of elements
  425.      */
  426.     template <class T>
  427.     void recv(T* arr, int size) {
  428.         MPI::Datatype data_type = get_type<T>();
  429.  
  430.         MPI::COMM_WORLD.Recv(
  431.             &arr[0], size, data_type, m_remote,
  432.             (m_message_tag == 0) ? MPI::ANY_TAG : m_message_tag, m_status);
  433.  
  434.         if (m_log_flag) {
  435.             tmp_log << "- receive array of type " << get_type_name<T>();
  436.             tmp_log << " of size=" << size << " from process=" << m_remote
  437.                     << endl;
  438.             flush();
  439.         }
  440.     }
  441.  
  442.     template <class T>
  443.     void recv(vector<T>& v) {
  444.         MPI::Datatype data_type = get_type<T>();
  445.  
  446.         int size = 0;
  447.         MPI::COMM_WORLD.Recv(
  448.             &size, 1, MPI::INT, m_remote,
  449.             (m_message_tag == 0) ? MPI::ANY_TAG : m_message_tag, m_status);
  450.  
  451.         v.resize(size);
  452.  
  453.         T* ptr = v.data();
  454.         MPI::COMM_WORLD.Recv(
  455.             ptr, size, data_type, m_remote,
  456.             (m_message_tag == 0) ? MPI::ANY_TAG : m_message_tag, m_status);
  457.  
  458.         if (m_log_flag) {
  459.             tmp_log << "- receive " << size << " elements of vector";
  460.             tmp_log << " of type " << get_type_name<T>();
  461.             tmp_log << " from process=" << m_remote << endl;
  462.             flush();
  463.         }
  464.     }
  465.  
  466.     /**
  467.      * Send array and receive value in return, this is an instance
  468.      * of the Sendrecv function.
  469.      * @param arr address of the array to send
  470.      * @param size size of the array to send
  471.      * @param value value to receive
  472.      */
  473.     template <class T, class U>
  474.     void sendrecv(T* array, int size, U& value) {
  475.         MPI::Datatype array_data_type = get_type<T>();
  476.         MPI::Datatype value_data_type = get_type<U>();
  477.  
  478.         if (m_log_flag) {
  479.             tmp_log << "- sendrecv/send array of type " << get_type_name<T>();
  480.             tmp_log << " of size=" << size << endl;
  481.             flush();
  482.         }
  483.  
  484.         MPI::COMM_WORLD.Sendrecv(&array[0], size, array_data_type, m_remote, 0,
  485.                                  &value, 1, value_data_type, MPI::ANY_SOURCE,
  486.                                  MPI::ANY_TAG, m_status);
  487.  
  488.         if (m_log_flag) {
  489.             tmp_log << "sendrecv/receive value=" << value << endl;
  490.             flush();
  491.         }
  492.     }
  493.  
  494.     /**
  495.      * Perform reduction
  496.      * @param item local data used to perform reduction
  497.      * @param total global data that will contain result
  498.      * @param op operation to perform (MPI::SUM, MPI::MAX, ...)
  499.      */
  500.     template <class T>
  501.     void reduce(T& item, T& total, const MPI::Op& op, int root_id = MASTER) {
  502.         MPI::Datatype data_type = get_type<T>();
  503.  
  504.         assert((root_id >= 0) and (root_id < m_nbr_processes));
  505.  
  506.         MPI::COMM_WORLD.Reduce(&item, &total, 1, data_type, op, root_id);
  507.  
  508.         if (m_log_flag) {
  509.             tmp_log << "- reduction gives value=" << total << endl;
  510.             flush();
  511.         }
  512.     }
  513.  
  514.     /**
  515.      * Perform gather operation
  516.      * @param lcl_array local array that is send to master process
  517.      * @param glb_array global array that will contain all local arrays
  518.      */
  519.     template <class T>
  520.     void gather(T* lcl_array, int size, T* glb_array, int root_id = MASTER) {
  521.         MPI::Datatype data_type = get_type<T>();
  522.  
  523.         assert((root_id >= 0) and (root_id < m_nbr_processes));
  524.  
  525.         MPI::COMM_WORLD.Gather(lcl_array, size, data_type, glb_array, size,
  526.                                data_type, root_id);
  527.  
  528.         if (m_log_flag) {
  529.             tmp_log << "- gather" << endl;
  530.             flush();
  531.         }
  532.     }
  533.  
  534.     /**
  535.      * Perform scatter operation
  536.      * @param big_array array of data that will be send by to all processors
  537.      * by parts
  538.      * @param big_size size of the array that will be sent
  539.      * @param small_array local array of data
  540.      */
  541.     template <class T>
  542.     void scatter(T* big_array, int big_size, T* small_array) {
  543.         MPI::Datatype data_type = get_type<T>();
  544.  
  545.         int small_size = big_size / m_nbr_processes;
  546.         MPI::COMM_WORLD.Scatter(big_array, small_size, data_type, small_array,
  547.                                 small_size, data_type, 0);
  548.  
  549.         if (m_log_flag) {
  550.             tmp_log << "- scatter array of size " << big_size;
  551.             tmp_log << " into " << m_nbr_processes << " chunks of "
  552.                     << small_size;
  553.             tmp_log << " elements" << endl;
  554.             flush();
  555.         }
  556.     }
  557.  
  558.     /**
  559.      * Perform scatter operation
  560.      * @param big_vec vector of data that will be send by to all processors
  561.      * by parts
  562.      * @param local_vec local vector of data
  563.      */
  564.     template <class T>
  565.     void scatter(vector<T>& big_vec, int big_size, vector<T>& local_vec,
  566.                  int begin = 0, int end = -1) {
  567.         MPI::Datatype data_type = get_type<T>();
  568.  
  569.         assert((begin >= 0) and (begin < big_size));
  570.         if (end == -1) end = big_size;
  571.         assert((begin < end) and (end <= big_size));
  572.         big_size = (end - begin + 1);
  573.  
  574.         int local_vec_size = big_size / m_nbr_processes;
  575.         MPI::COMM_WORLD.Scatter(big_vec.data(), local_vec_size, data_type,
  576.                                 local_vec.data(), local_vec_size, data_type, 0);
  577.  
  578.         if (m_log_flag) {
  579.             tmp_log << "- scatter vector of size " << big_size;
  580.             tmp_log << " into " << m_nbr_processes << " chunks of "
  581.                     << local_vec_size;
  582.             tmp_log << " elements" << endl;
  583.             flush();
  584.         }
  585.     }
  586.  
  587.     // MPI_Bcast(void *buffer, int count, MPI_Datatype datatype, int root,
  588.     // MPI_Comm comm)
  589.     template <class T>
  590.     void broadcast(T* glb_array, int size, int root_id = MASTER) {
  591.         MPI::Datatype data_type = get_type<T>();
  592.  
  593.         assert((root_id >= 0) and (root_id < m_nbr_processes));
  594.  
  595.         MPI::COMM_WORLD.Bcast((void*)glb_array, size, data_type, root_id);
  596.  
  597.         if (m_log_flag) {
  598.             tmp_log << "- broadcast array of size " << size << endl;
  599.             flush();
  600.         }
  601.     }
  602.  
  603.     template <class T>
  604.     void broadcast(std::vector<T>& vec, int root_id = MASTER) {
  605.         MPI::Datatype data_type = get_type<T>();
  606.  
  607.         assert((root_id >= 0) and (root_id < m_nbr_processes));
  608.  
  609.         MPI::COMM_WORLD.Bcast((void*)vec.data(), static_cast<int>(vec.size()),
  610.                               data_type, root_id);
  611.  
  612.         if (m_log_flag) {
  613.             tmp_log << "- broadcast array of size " << vec.size()
  614.                     << " elements";
  615.             tmp_log << " from process=" << root_id << endl;
  616.             flush();
  617.         }
  618.     }
  619.  
  620.     typedef std::ostream& (*ManipFn)(std::ostream&);
  621.     typedef std::ios_base& (*FlagsFn)(std::ios_base&);
  622.  
  623.     template <class T>
  624.     Process& operator<<(vector<T>& v) {
  625.         tmp_log << '[';
  626.         if (v.size() > 0) {
  627.             tmp_log << v[0];
  628.             for (int i = 1; i < static_cast<int>(v.size()); ++i) {
  629.                 tmp_log << ' ' << v[i];
  630.             }
  631.         }
  632.         tmp_log << ']';
  633.         return *this;
  634.     }
  635.  
  636.     template <class T>  // int, double, strings, etc
  637.     Process& operator<<(const T& output) {
  638.         tmp_log << output;
  639.         return *this;
  640.     }
  641.  
  642.     // endl, flush, setw, setfill, etc.
  643.     Process& operator<<(ManipFn manip) {
  644.         manip(tmp_log);
  645.         if (manip == static_cast<ManipFn>(std::flush) ||
  646.             manip == static_cast<ManipFn>(std::endl)) {
  647.             this->flush();
  648.         }
  649.         return *this;
  650.     }
  651.  
  652.     // setiosflags, resetiosflags
  653.     Process& operator<<(FlagsFn manip) {
  654.         manip(tmp_log);
  655.         return *this;
  656.     }
  657.     void logs(ostream& out) {
  658.         MPI::COMM_WORLD.Barrier();
  659.         m_verbose_flag = m_log_flag = false;
  660.         if (m_id == 0) {
  661.             out.flush();
  662.             out << std::endl;
  663.             out << "=====================" << std::endl;
  664.             out << "===    LOGGING    ===" << std::endl;
  665.             out << "=====================" << std::endl;
  666.             out << "---------------------" << std::endl;
  667.             out << "CPU " << m_id << " (pid=" << unix_pid() << ")" << std::endl;
  668.             out << "---------------------" << std::endl;
  669.             out << log_stream.str();
  670.             out.flush();
  671.             remote(1);
  672.             int token = -255;
  673.             send(token);
  674.         } else {
  675.             remote(m_id - 1);
  676.             int token;
  677.             recv(token);
  678.             out << "---------------------" << std::endl;
  679.             out << "CPU " << m_id << " (pid=" << unix_pid() << ")" << std::endl;
  680.             out << "---------------------" << std::endl;
  681.             out << log_stream.str();
  682.             out.flush();
  683.             if (m_id < m_nbr_processes - 1) {
  684.                 remote(m_id + 1);
  685.                 token = -255;
  686.                 send(token);
  687.             }
  688.         }
  689.     }
  690.  
  691.     typedef void (*Code)(Process& p);
  692.  
  693.     void run(Code code) { code(*this); }
  694. };
  695.  
  696. }  // end of namespace mpi
  697.  
  698. }  // end of namespace ez
  699.  

On peut alors réécrire les programes précédents de manière plus simple :

Afficher le code    ens/m2/paracpu/ezmpi_send_and_recv_array.cpp
  1. // ==================================================================
  2. // Author: Jean-Michel Richer
  3. // Email: jean-michel.richer@univ-angers.fr
  4. // Date: Aug 2020
  5. // Last modified: September 2026
  6. // Purpose: Demonstrate basic send functionality, send array of
  7. // integers from master (rank=0) to slave (rank=1)
  8. // ==================================================================
  9. // Don't forget to launch with at least two processes:
  10. //
  11. // mpirun -n 2 bin/ezmpi_send_and_recv.exe
  12. // or
  13. // mpirun -n 4 bin/ezmpi_send_and_recv.exe
  14. //
  15. // ==================================================================
  16. #include <unistd.h>
  17.  
  18. #include <cstdlib>
  19. #include <iostream>
  20. #include <numeric>
  21. using namespace std;
  22. #include "ezmpi.hpp"
  23. using namespace ez::mpi;
  24.  
  25. /**
  26.  * @brief return a string representation of an array
  27.  */
  28. string to_string(int* array, int size) {
  29.     ostringstream oss;
  30.     oss << '[';
  31.     if (size > 0) {
  32.         oss << array[0];
  33.         for (int i = 1; i < size; ++i) {
  34.             oss << ',' << array[i];
  35.         }
  36.     }
  37.     oss << ']';
  38.     return oss.str();
  39. }
  40.  
  41. /**
  42.  * run master and slaves
  43.  */
  44. void run(int argc, char* argv[]) {
  45.     // data to send
  46.     int* array_data = nullptr;
  47.     int array_size;
  48.  
  49.     // create process here so that it will be destroyed
  50.     // at the end of the procedure
  51.     Process process(argc, argv);
  52.     process.log(true);
  53.  
  54.     if (process.is_master()) {
  55.         // ----------------------------------------------------------
  56.         // Code for the master
  57.         // ----------------------------------------------------------
  58.  
  59.         // set remote cpu identifier that will communicate with master
  60.         process.remote(1);
  61.  
  62.         array_size = 10;
  63.         array_data = new int[array_size];
  64.         iota(&array_data[0], &array_data[array_size], 1);
  65.  
  66.         process << to_string(array_data, array_size) << endl;
  67.  
  68.         // send size of the array to slave1
  69.         process.send(array_size);
  70.         // send data of array to slave1
  71.         process.send(array_data, array_size);
  72.  
  73.         // wait for slave1 to compute and send back sum
  74.         float result = 0;
  75.         process.recv(result);
  76.  
  77.         process << "- result is " << result << " for length=" << array_size;
  78.         process << ", expected=" << (array_size * (array_size + 1)) / 2 << endl;
  79.  
  80.     } else if (process.is_slave(1)) {
  81.         // ----------------------------------------------------------
  82.         // Code for the first slave
  83.         // ----------------------------------------------------------
  84.  
  85.         // tells to Process to send data to master processor
  86.         process.remote(0);
  87.  
  88.         // process of id 1 will receive the array and compute the sum
  89.         // we receive the array length and allocate space
  90.         process.recv(array_size);
  91.  
  92.         array_data = new int[array_size];
  93.  
  94.         // we receive the data
  95.         process.recv(array_data, array_size);
  96.  
  97.         process << "- receive array " << to_string(array_data, array_size)
  98.                 << endl;
  99.  
  100.         // 3- compute sum
  101.         float sum = std::accumulate(&array_data[0], &array_data[array_size], 0);
  102.  
  103.         // 4- send back result
  104.         process.send(sum);
  105.  
  106.     } else {
  107.         // ----------------------------------------------------------
  108.         // Code for the other processes: don't do anything
  109.         // ----------------------------------------------------------
  110.  
  111.         process << "- is idle" << endl;
  112.     }
  113.  
  114.     process.logs(cout);
  115. }
  116.  
  117. /**
  118.  * main function
  119.  *
  120.  */
  121. int main(int argc, char** argv) {
  122.     srand(time(nullptr));
  123.  
  124.     // call a function that contains the code for the master and
  125.     // the slaves in order to avoid problems of initialization and
  126.     // finalization of processes
  127.     run(argc, argv);
  128.  
  129.     exit(EXIT_SUCCESS);
  130. }
Afficher le code    ens/m2/paracpu/ezmpi_send_and_recv_vector.cpp
  1. // ==================================================================
  2. // Author: Jean-Michel Richer
  3. // Email: jean-michel.richer@univ-angers.fr
  4. // Date: Aug 2020
  5. // Last modified: September 2026
  6. // Purpose: Demonstrate basic send functionality, send array of float
  7. // from master (rank=0) to slave (rank=1)
  8. // ==================================================================
  9. // Don't forget to launch with
  10. //
  11. // mpirun -n 2 bin/ezmpi_send_and_recv.exe
  12. //
  13. // ==================================================================
  14. #include <unistd.h>
  15. #include <unistd.h>  // for sleep
  16.  
  17. #include <cstdlib>
  18. #include <iostream>
  19. #include <numeric>
  20. using namespace std;
  21. #include "ezmpi.hpp"
  22. using namespace ez::mpi;
  23.  
  24. /**
  25.  * run master and slaves
  26.  */
  27. void run(int argc, char* argv[]) {
  28.     // data to send
  29.     vector<int> vec;
  30.     int vec_size;
  31.  
  32.     // create process here so that it will be destroyed
  33.     // at the end of the procedure
  34.     Process process(argc, argv);
  35.     process.log(true);
  36.  
  37.     if (process.is_master()) {
  38.         // ----------------------------------------------------------
  39.         // Code for the master
  40.         // ----------------------------------------------------------
  41.  
  42.         // set remote cpu identifier that will communicate with master
  43.         process.remote(1);
  44.  
  45.         vec_size = 10;
  46.         vec.resize(vec_size);
  47.         iota(vec.begin(), vec.end(), 1);
  48.  
  49.         process << vec << endl;
  50.  
  51.         // send vector to slave1
  52.         process.send(vec);
  53.  
  54.         // wait for slave1 to compute and send back sum
  55.         float result = 0;
  56.         process.recv(result);
  57.  
  58.         process << "- result is " << result << " for length=" << vec_size;
  59.         process << ", expected=" << (vec_size * (vec_size + 1)) / 2 << endl;
  60.  
  61.     } else if (process.is_slave(1)) {
  62.         // ----------------------------------------------------------
  63.         // Code for the first slave
  64.         // ----------------------------------------------------------
  65.  
  66.         // tells to Process to send data to master processor
  67.         process.remote(0);
  68.  
  69.         // process of id 1 will receive the vector and compute the sum
  70.  
  71.         // we receive the vector
  72.         process.recv(vec);
  73.  
  74.         process << "- receive vector " << vec << endl;
  75.  
  76.         // 3- compute sum
  77.         float sum = std::accumulate(vec.begin(), vec.end(), 0);
  78.  
  79.         // 4- send back result
  80.         process.send(sum);
  81.  
  82.     } else {
  83.         // ----------------------------------------------------------
  84.         // Code for the other processes: don't do anything
  85.         // ----------------------------------------------------------
  86.  
  87.         process << "- is idle" << endl;
  88.     }
  89.  
  90.     process.logs(cout);
  91. }
  92.  
  93. /**
  94.  * main function
  95.  *
  96.  */
  97. int main(int argc, char** argv) {
  98.     srand(time(nullptr));
  99.  
  100.     // call a function that contains the code for the master and
  101.     // the slaves in order to avoid problems of initialization and
  102.     // finalization of processes
  103.     run(argc, argv);
  104.  
  105.     exit(EXIT_SUCCESS);
  106. }
Afficher le code    ens/m2/paracpu/ezmpi_sendrecv.cpp
  1. // ==================================================================
  2. // Author: Jean-Michel Richer
  3. // Email: jean-michel.richer@univ-angers.fr
  4. // Date: Aug 2020
  5. // Last modified: September 2026
  6. // Purpose: Demonstrate basic send functionality, send array of float
  7. // from master (rank=0) to slave (rank=1)
  8. // ==================================================================
  9. #include <unistd.h>
  10. #include <unistd.h>  // for sleep
  11.  
  12. #include <cstdlib>
  13. #include <iostream>
  14. #include <numeric>
  15. #include <sstream>
  16. using namespace std;
  17. #include "ezmpi.hpp"
  18. using namespace ez::mpi;
  19.  
  20. /**
  21.  * run master and slaves
  22.  * - the master:
  23.  */
  24. void run(int argc, char* argv[]) {
  25.     // data to send
  26.     vector<int> vec;
  27.     int vec_size;
  28.  
  29.     // create process here so that it will be destroyed
  30.     // at the end of the procedure
  31.     Process process(argc, argv);
  32.  
  33.     if (process.is_master()) {
  34.         // ----------------------------------------------------------
  35.         // Code of the master
  36.         // ----------------------------------------------------------
  37.  
  38.         // set remote cpu identifier that will communicate with master
  39.         process.remote(1);
  40.  
  41.         // processor of id 0 (master) fills and sends the array
  42.         vec_size = 10;
  43.         vec.resize(vec_size);
  44.         iota(vec.begin(), vec.end(), 1);
  45.  
  46.         // 1- send size of array to remote
  47.         process.send(vec_size);
  48.  
  49.         // 2- send data of array to remote and receive result
  50.         float result = 0;
  51.         process.sendrecv(vec.data(), vec_size, result);
  52.  
  53.         process << "- result is " << result << " for length=" << vec_size;
  54.         process << ", expected=" << (vec_size * (vec_size + 1)) / 2 << endl;
  55.  
  56.     } else if (process.id() == ez::mpi::SLAVE_1) {
  57.         // ----------------------------------------------------------
  58.         // Code for the first slave
  59.         // ----------------------------------------------------------
  60.  
  61.         // tells to Process to send data to master processor
  62.         process.remote(0);
  63.  
  64.         // process of id 1 will receive the array and compute the sum
  65.         // 1- we receive the array length and allocate space
  66.         process.recv(vec_size);
  67.  
  68.         vec.resize(vec_size);
  69.  
  70.         // 2- we receive the data
  71.         process.recv(vec.data(), vec_size);
  72.  
  73.         // 3- compute sum
  74.         float sum = accumulate(vec.begin(), vec.end(), 0);
  75.  
  76.         // 4- send back result
  77.         process.send(sum);
  78.  
  79.     } else {
  80.         // ----------------------------------------------------------
  81.         // other processes (if any): don't do anything
  82.         // ----------------------------------------------------------
  83.         process << "- is idle" << endl;
  84.     }
  85.  
  86.     process.logs(cout);
  87. }
  88.  
  89. // ==============================================================
  90. // version C++
  91. // ==============================================================
  92. int main(int argc, char** argv) {
  93.     // call a function that contains the code for the master and
  94.     // the slaves in order to avoid problems of initialization and
  95.     // finalization of processes
  96.     run(argc, argv);
  97.  
  98.     exit(EXIT_SUCCESS);
  99. }

5.7. Réduction

Certains traitements comme la réduction ou le scan sont implantés spécifiquement pour MPI.

Par exemple pour la réduction :

Afficher le code    ens/m2/paracpu/ezmpi_reduce_sum.cpp
  1. // ==================================================================
  2. // Author: Jean-Michel Richer
  3. // Email: jean-michel.richer@univ-angers.fr
  4. // Date: Aug 2020
  5. // Last modified: September 2026
  6. // Purpose: Demonstrate basic reduce functionality by computing the
  7. // sum of a floatting point values on each process
  8. // ==================================================================
  9. #include <unistd.h>
  10. #include <unistd.h>  // for sleep
  11.  
  12. #include <cstdlib>
  13. #include <iostream>
  14. #include <sstream>
  15. using namespace std;
  16. #include "./ezmpi.hpp"
  17. using namespace ez::mpi;
  18.  
  19. /**
  20.  * run master and slaves
  21.  */
  22. void run(int argc, char* argv[]) {
  23.     Process process(argc, argv);
  24.     process.verbose(false);
  25.  
  26.     float local_data = process.id() + 1;
  27.  
  28.     process << "- local data=" << local_data << endl;
  29.  
  30.     float reduction = 0;
  31.  
  32.     // reduction must be performed outside the master and slaves
  33.     process.reduce(local_data, reduction, MPI::SUM);
  34.  
  35.     // only the master will report the result
  36.     if (process.is_master()) {
  37.         cout << "reduction = " << reduction << endl;
  38.     }
  39.  
  40.     process.logs(cout);
  41. }
  42.  
  43. // ==============================================================
  44. // version C++
  45. // ==============================================================
  46. int main(int argc, char** argv) {
  47.     // call a function that contains the code for the master and
  48.     // the slaves in order to avoid problems of initialization and
  49.     // finalization of processes
  50.     run(argc, argv);
  51.  
  52.     exit(EXIT_SUCCESS);
  53. }
Afficher le code    ens/m2/paracpu/ezmpi_reduce_min.cpp
  1. // ==================================================================
  2. // Author: Jean-Michel Richer
  3. // Email: jean-michel.richer@univ-angers.fr
  4. // Date: Aug 2020
  5. // Last modified: September 2026
  6. // Purpose: Demonstrate basic minimum functionaly on an array
  7. // The master create a big array and sends part of it to the
  8. // slaves
  9. // ==================================================================
  10. #include <unistd.h>  // for sleep
  11.  
  12. #include <algorithm>
  13. #include <cstdlib>
  14. #include <iostream>
  15. #include <random>
  16. #include <sstream>
  17. #include <vector>
  18. using namespace std;
  19. #include "./ezmpi.hpp"
  20. using namespace ez::mpi;
  21. #include <getopt.h>
  22.  
  23. int initialization_method = 1;
  24.  
  25. /**
  26.  * run master and slaves
  27.  */
  28. void run(int argc, char* argv[]) {
  29.     Process process(argc, argv);
  30.     process.verbose(false);
  31.  
  32.     // the big vector for the master
  33.     int big_vec_size = 100;
  34.     vector<int> big_vec;
  35.  
  36.     // reserve space for the local vector on eaxh process
  37.     int local_vec_size = big_vec_size / process.nbr_processes();
  38.     vector<int> local_vec(local_vec_size);
  39.  
  40.     if (process.is_master()) {
  41.         big_vec.resize(big_vec_size);
  42.  
  43.         process << "- master local vector size=" << local_vec_size << endl;
  44.  
  45.         // obtain a random seed from the hardware
  46.         random_device rd;
  47.         // standard Mersenne Twister engine seeded with rd()
  48.         mt19937 gen(rd());
  49.  
  50.         uniform_int_distribution<> distrib(0, big_vec_size * 2);
  51.  
  52.         if (initialization_method == 1) {
  53.             iota(big_vec.begin(), big_vec.end(), 1);
  54.             reverse(big_vec.begin(), big_vec.end());
  55.         } else {
  56.             // generate random values for the big vector
  57.             generate(big_vec.begin(), big_vec.end(),
  58.                      [&]() { return distrib(gen); });
  59.         }
  60.  
  61.         process << "- big vector=" << big_vec << endl;
  62.     }
  63.  
  64.     // distribution from the big vector to the local vectors as slices
  65.     // for example if big_vec_size=100 et the number of processes is 4
  66.     // each process will get 25 elements
  67.     process.scatter(big_vec, big_vec_size, local_vec);
  68.  
  69.     process << "- local vector=" << local_vec << endl;
  70.  
  71.     // each process (included the master) finds the minimum value in its
  72.     // local vector
  73.     int local_minimum = *std::min_element(local_vec.begin(), local_vec.end());
  74.  
  75.     // each process reports the local minimum found in its local vector
  76.     process << "- local minimum=" << local_minimum << endl;
  77.  
  78.     // perform reduction of a minimum
  79.     // note: this code must be excuted by all processes
  80.     int global_minimum = -1;
  81.     process.reduce(local_minimum, global_minimum, MPI::MIN);
  82.  
  83.     if (process.is_master()) {
  84.         process << "- global minimum=" << global_minimum << endl;
  85.     }
  86.  
  87.     process.logs(cout);
  88. }
  89.  
  90. // ==============================================================
  91. // version C++
  92. // ==============================================================
  93. int main(int argc, char** argv) {
  94.     int opt;
  95.     while ((opt = getopt(argc, argv, "i:")) != -1) {
  96.         switch (opt) {
  97.             case 'i':
  98.                 initialization_method = std::stoi(optarg);
  99.                 break;
  100.             default:
  101.                 cerr << "use -i 1 or -i 2 to initialize the big vector" << endl;
  102.         }
  103.     }
  104.  
  105.     if ((initialization_method < 1) or (initialization_method > 2)) {
  106.         initialization_method = 1;
  107.     }
  108.  
  109.     // call a function that contains the code for the master and
  110.     // the slaves in order to avoid problems of initialization and
  111.     // finalization of processes
  112.     run(argc, argv);
  113.  
  114.     exit(EXIT_SUCCESS);
  115. }

5.8. Gather et scatter

MPI propose deux opérations inverses :

  • scatter (split) qui permet de décomposer un tableau en différentes parties égales et de les envoyer aux autres processus
  • gather (join) qui exécute l'opération inverse : à partir de tableaux locaux détenus par chacun des processus, on crée un tableau global plus grand qui est la concaténation des tableaux locaux

scatter gather

Voici un exemple pour lequel $K$ processus créent des tableaux de 10 entiers qui sont envoyés au master (processus de rang 0) :

  • Processus 0 : [1, 2, ... , 10]
  • Processus 1 : [11, 12, ..., 20]
  • Processus $i$ : [$(i*10+1)$, ...,$(i+1)*10$ ]
Afficher le code    ens/m2/paracpu/ezmpi_scatter.cpp
  1. // ==================================================================
  2. // Author: Jean-Michel Richer
  3. // Date: 20 Aug 2016
  4. // Purpose: Use of MPI::Scatter
  5. // The master process creates its own local data and they are send
  6. // to other process (of rank not equal to 0)
  7. // ==================================================================
  8. #include <unistd.h>
  9. #include <unistd.h>  // for sleep
  10.  
  11. #include <cstdlib>
  12. #include <iostream>
  13. #include <numeric>
  14. using namespace std;
  15. #include "ezmpi.hpp"
  16. using namespace ez::mpi;
  17.  
  18. /**
  19.  * run master and slaves
  20.  */
  21. void run(int argc, char* argv[]) {
  22.     Process process(argc, argv);
  23.  
  24.     // local vector
  25.     int local_vec_size = 10;
  26.     vector<int> local_vec(local_vec_size, 0);
  27.  
  28.     // big vector
  29.     int big_vec_size = process.nbr_processes() * local_vec_size;
  30.     vector<int> big_vec(big_vec_size);
  31.  
  32.     if (process.is_master()) {
  33.         // master process will send data to others
  34.  
  35.         iota(big_vec.begin(), big_vec.end(), 1);
  36.  
  37.         process << "- big vector = " << big_vec << endl;
  38.     }
  39.  
  40.     // Master sends data to all others
  41.     process.scatter(big_vec, big_vec_size, local_vec);
  42.  
  43.     // Each process reports the data it has received
  44.     process << "- local vector = " << local_vec << endl;
  45.  
  46.     process.logs(cout);
  47. }
  48.  
  49. // ==================================================================
  50. // C++ version
  51. // ==================================================================
  52. int main(int argc, char** argv) {
  53.     // call a function that contains the code for the master and
  54.     // the slaves in order to avoid problems of initialization and
  55.     // finalization of processes
  56.     run(argc, argv);
  57.  
  58.     exit(EXIT_SUCCESS);
  59. }
> mpirun -n 4 ./ezmpi_scatter.exe
====================
=== FINAL RESULT ===
====================
---------------------
CPU 0
---------------------
14:15:05 cpu 0/4: global array = [1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 ]
14:15:05 cpu 0/4: scatter
14:15:05 cpu 0/4: local array=[1 2 3 4 5 6 7 8 9 10 ]
---------------------
CPU 1
---------------------
14:15:05 cpu 1/4: scatter
14:15:05 cpu 1/4: local array=[11 12 13 14 15 16 17 18 19 20 ]
---------------------
CPU 2
---------------------
14:15:05 cpu 2/4: scatter
14:15:05 cpu 2/4: local array=[21 22 23 24 25 26 27 28 29 30 ]
---------------------
CPU 3
---------------------
14:15:05 cpu 3/4: scatter
14:15:05 cpu 3/4: local array=[31 32 33 34 35 36 37 38 39 40 ]
Afficher le code    ens/m2/paracpu/ezmpi_gather.cpp
  1. // ==================================================================
  2. // Author: Jean-Michel Richer
  3. // Date: Aug 2020
  4. // Last modified: November 2020
  5. // Purpose: Use of MPI::Gather
  6. // Each process creates its own local vector of data and they are all
  7. // sent to the master process (of rank 0)
  8. // ==================================================================
  9. #include <unistd.h>
  10. #include <unistd.h>  // for sleep
  11.  
  12. #include <cstdlib>
  13. #include <iostream>
  14. #include <numeric>
  15. #include <sstream>
  16. #include <vector>
  17. using namespace std;
  18. #include "ezmpi.hpp"
  19. using namespace ez::mpi;
  20.  
  21. /**
  22.  * run master and slaves
  23.  */
  24. void run(int argc, char* argv[]) {
  25.     Process process(argc, argv);
  26.  
  27.     // --------------------------------------------------------------
  28.     // declaration of data
  29.     // --------------------------------------------------------------
  30.  
  31.     // this is for the master
  32.     vector<int> big_vec;
  33.  
  34.     // this is for the master and the slaves
  35.     int local_vec_size = 10;
  36.     vector<int> local_vec(local_vec_size);
  37.  
  38.     // depending on the process fill the local vector with numbers
  39.     iota(local_vec.begin(), local_vec.end(), process.id() * 10 + 1);
  40.  
  41.     process << "- local vector=" << local_vec << endl;
  42.  
  43.     if (process.is_master()) {
  44.         // master process will gather data from others
  45.         // so define the size of the vector that will
  46.         // receive other data from the slaves and
  47.         // fill it with 0's
  48.         big_vec.resize(process.nbr_processes() * local_vec_size, 0);
  49.     }
  50.  
  51.     // gather data from all processes to master
  52.     // note: this must be executed by all processes
  53.     process.gather(local_vec.data(), local_vec_size, big_vec.data());
  54.  
  55.     // finally print the content of the big vector on the master
  56.     if (process.is_master()) {
  57.         process << "- big vector = " << big_vec << endl;
  58.     }
  59.  
  60.     process.logs(cout);
  61. }
  62.  
  63. // ==================================================================
  64. // C++ version
  65. // ==================================================================
  66. int main(int argc, char** argv) {
  67.     // call a function that contains the code for the master and
  68.     // the slaves in order to avoid problems of initialization and
  69.     // finalization of processes
  70.     run(argc, argv);
  71.  
  72.     exit(EXIT_SUCCESS);
  73. }
> mpirun -n 4 ./ezmpi_gather.exe
====================
=== FINAL RESULT ===
====================
---------------------
CPU 0
---------------------
14:19:02 cpu 0/4: local array=[1 2 3 4 5 6 7 8 9 10 ]
14:19:02 cpu 0/4: gather
14:19:02 cpu 0/4: global array = [1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 ]
---------------------
CPU 1
---------------------
14:19:02 cpu 1/4: local array=[11 12 13 14 15 16 17 18 19 20 ]
14:19:02 cpu 1/4: gather
---------------------
CPU 2
---------------------
14:19:02 cpu 2/4: local array=[21 22 23 24 25 26 27 28 29 30 ]
14:19:02 cpu 2/4: gather
---------------------
CPU 3
---------------------
14:19:02 cpu 3/4: local array=[31 32 33 34 35 36 37 38 39 40 ]
14:19:02 cpu 3/4: gather

5.9. Broadcast (diffuser)

Enfin, l'opération broadcast envoie à tous les esclaves une copie d'un tableau détenu par le maître par exemple :

Afficher le code    ens/m2/paracpu/ezmpi_broadcast.cpp
  1. // ==================================================================
  2. // Author: Jean-Michel Richer
  3. // Date: Aug 2020
  4. // Last modified: November 2020
  5. // Purpose: Use of MPI::BCast
  6. // The master process creates its own local data and they are send
  7. // to other processes (that have a rank not equal to the master)
  8. // ==================================================================
  9. #include <unistd.h>
  10. #include <unistd.h>  // for sleep
  11.  
  12. #include <cstdlib>
  13. #include <iostream>
  14. #include <numeric>
  15. using namespace std;
  16. #include "ezmpi.hpp"
  17. using namespace ez::mpi;
  18.  
  19. /**
  20.  * run master and slaves
  21.  */
  22. void run(int argc, char* argv[]) {
  23.     Process process(argc, argv);
  24.  
  25.     // declare a vector of vec_size integers
  26.     int vec_size = 10;
  27.     vector<int> vec(vec_size, 0);
  28.  
  29.     if (process.is_master()) {
  30.         // master process will:
  31.         // 1- fille the vector
  32.         // 2- send the vector to other processes (see broadcast below)
  33.         iota(vec.begin(), vec.end(), 1);
  34.  
  35.         process << "- create global array =" << vec << endl;
  36.     }
  37.  
  38.     // the master process sends the vector to all other processes
  39.     // the instruction must be executed by all other processes
  40.     process.broadcast(vec);
  41.  
  42.     // Each process reports the data it has received
  43.     if (process.id() != MASTER) {
  44.         process << "- local copy array=" << vec << endl;
  45.     }
  46.  
  47.     process.logs(cout);
  48. }
  49.  
  50. // ==================================================================
  51. // C++ version
  52. // ==================================================================
  53. int main(int argc, char** argv) {
  54.     // call a function that contains the code for the master and
  55.     // the slaves in order to avoid problems of initialization and
  56.     // finalization of processes
  57.     run(argc, argv);
  58.  
  59.     exit(EXIT_SUCCESS);
  60. }
> mpirun -n 4 ./ezmpi_broadcast.exe
====================
=== FINAL RESULT ===
====================
---------------------
CPU 0
---------------------
17:08:27 cpu 0/4: global array on master = [1 2 3 4 5 6 7 8 9 10 ]
17:08:27 cpu 0/4: broadcast
17:08:27 cpu 0/4: global array (after broadcast) =[1 2 3 4 5 6 7 8 9 10 ]
---------------------
CPU 1
---------------------
17:08:27 cpu 1/4: global array on slave = [0 0 0 0 0 0 0 0 0 0 ]
17:08:27 cpu 1/4: broadcast
17:08:27 cpu 1/4: global array (after broadcast) =[1 2 3 4 5 6 7 8 9 10 ]
---------------------
CPU 2
---------------------
17:08:27 cpu 2/4: global array on slave = [0 0 0 0 0 0 0 0 0 0 ]
17:08:27 cpu 2/4: broadcast
17:08:27 cpu 2/4: global array (after broadcast) =[1 2 3 4 5 6 7 8 9 10 ]
---------------------
CPU 3
---------------------
17:08:27 cpu 3/4: global array on slave = [0 0 0 0 0 0 0 0 0 0 ]
17:08:27 cpu 3/4: broadcast
17:08:27 cpu 3/4: global array (after broadcast) =[1 2 3 4 5 6 7 8 9 10 ]

5.10. Execution sur plusieurs machines

Exécuter un programme MPI sur plusieurs machines relève de la gageure, j'ai passé 4 heures pour pouvoir y parvenir après moult problèmes :

  • version d'openmpi déjà installée sur la machine secondaire (cf. ci-après AMD Ryzen 5 5600G) qui créait des conflits et empêchait toute exécution sur la machine secondaire
  • installation d'une carte graphique sur la machine secondaire, pas une vieille carte qui n'est pas reconnue, une récente ;-), mais casse d'une nappe SATA et d'une antenne Wifi lors du remplacement de la carte graphique et du changement de disque dur pour réinstallation d'Ubuntu
  • réinstallation totale de la machine secondaire : OS, compilateur g++, autoconf, automake, libtools, make, openssh-server, openmpi-4.1.8
  • réinstallation d'openmpi-4.1.8 sur la machine principale dans /usr/local
  • création de clé ssh qui apparaissait plusieurs fois empêchant la connexion ssh
  • longues conversations avec ChatGPT pour tenter de résoudre les problèmes pas à pas

Pour que cela fonctionne il faut, sur toutes les machines :

  • disposer du même système d'exploitation (ou de versions proches)
  • avoir la même version d'OpenMPI (ex : 4.1.8) installée dans les mêmes répertoires (ex : /usr/local/openmpi-4.1.8)
  • disposer des mêmes compilateurs C++
  • avoir la même hierarchie de répertoires, là où se trouvent vos fichiers MPI (ex : ~/dev/cpp/parallelism)
  • compiler le même programme (ou une version proche)
  • créer une clé ssh et la partager avec toutes les machines (secondaires) qui exécuteront le code pour pouvoir se connecter en SSH automatiquement sans avoir à saisir de mot de passe

Dans l'exemple qui suit, on dispose de deux machines :

  • solaris : 192.168.1.194 (machine principale, AMD Ryzen 5 9600X)
  • jupiter : 192.168.1.109 (machine secondaire, AMD Ryzen 5 5600G)

Eventuellement, modifiez /etc/hosts en ajoutant les lignes suivantes pour pouvoir faire référence aux machines par leur nom et non leur IP :

192.168.1.194 solaris
192.168.1.109 jupiter

5.10.1. Génération de la clé ssh

Sur solaris :

solaris> ssh-keygen -t ed25519 -N "" -f ~/.ssh/id_ed25519
solaris> ssh-copy-id richer@192.168.1.109

Se connecter une première fois afin d'activer la connexion ssh et la clé :

solaris> ssh richer@192.168.1.109
Welcome to Ubuntu 24.04 LTS (GNU/Linux 6.8.0-31-generic x86_64)
...
jupiter> exit

5.10.2. Compilation

Sur chaque machine compiler le programme à exécuter :

> mpic++ -o equation_mpi.exe equation_mpi.cpp -O3

On peut probablement le compiler sur la machine principale et le copier sur les machines secondaires par scp si on dispose des mêmes OS et compilateurs.

5.10.3. Fichier hosts.txt

Ce fichier décrit le nombre de slots (threads / processus) sur chaque machine, il faut le créer sur la machine principale :

solaris@~dev/cpp/parallelism> cat hosts.txt
solaris slots=12
jupiter slots=12

5.10.4. Exécution du programme

Pour exécuter le programme avec 16 processus, il faut procéder ainsi sur la machine principale :

solaris@~/dev/cpp/parallelism> time mpirun -np 16 --hostfile hosts.txt \
   --host solaris:8,jupiter:8  --prefix /usr/local/openmpi-4.1.8 --map-by ppr:8:node  \
   --mca btl tcp,self --mca btl_tcp_if_include 192.168.1.0/24 \
   --mca oob_tcp_if_include 192.168.1.0/24 \
   ~/dev/cpp/parallelism/equation_mpi.exe -n 26
2 0 0 0 0 4 0   # 1*2 + 6*4 = 26
4 0 0 0 2 2 0   # 1*4 + 5*2 + 6*2 = 26
6 0 0 0 0 1 2   # ...
8 0 0 0 0 3 0 
10 0 0 0 2 1 0 
12 0 0 0 0 0 2 
14 0 0 0 0 2 0 
16 0 0 0 2 0 0 
18 0 0 2 0 0 0 
20 0 0 0 0 1 0 
22 0 0 1 0 0 0 
23 0 1 0 0 0 0 
0 0 0 0 0 2 2 
26 0 0 0 0 0 0 
24 1 0 0 0 0 0 

real    0m6,095s
user    0m18,259s
sys     0m2,002s
  • on déclare utiliser 16 programmes : -np 16
  • dont 8 sur solaris et 8 sur jupiter : --host solaris:8,jupiter:8
  • la politique de répartition des programmes par ressource (Processes Per Resource) est de 8 sur chaque machine : --map-by ppr:8:node
  • les options --mca ... ne sont pas forcément nécessaires

5.11. Liens

5.12. Exercices

Exercice 5.1

Une suite de Syracuse est une suite d'entiers naturels définie de la manière suivante : on part d'un nombre entier strictement positif

  • s'il est pair, on le divise par 2
  • s'il est impair, on le multiplie par 3 et on ajoute 1

Cette suite possède la propriété de converger vers 1 après un certain nombre d'étapes.

Mettre en place une solution MPI avec trois instances :

  • le maître génére un nombre entier aléatoire $x$
  • si le nombre $x$ est pair, il l'envoie à l'esclave n°1 qui retourne $x/2$
  • si le nombre $x$ est impair, il l'envoie à l'esclave n°2 qui retourne $3x+1$
  • à chaque étape, le maître affiche la nouvelle valeur de $x$

Le maître comptera le nombre d'appels à l'esclave n°1 et l'esclave n°2.

Lorsque le maître recoît la valeur 1, il s'arrête et envoie un code d'arrêt aux esclaves (un nombre négatif par exemple). Puis il affiche le nombre d'appels à chaque esclave.

Exercice 5.2

Objectif : Écrire un programme MPI pour rechercher un élément dans un tableau.

Description :

  • créer 5 processus
  • générer un tableau d'un million d'entiers aléatoires sur le processus maître (processus de rang 0)
  • répartir le tableau entre les différents processus esclaves en utilisant la fonction MPI_Scatter
  • le maître génère un nombre aléatoire $x$ et demande à chaque esclave de rechercher le nombre d'occurrences de $x$ dans son tableau et de retourner le résultat
  • afficher sur le maître, le nombre d'occurrences par esclave