Reading Time: 10 minutes

Points à retenir clés

  • Le code Python peut s’exécuter sur un ordinateur portable et un système HPC volumineux avec la même logique de base lorsque le flux de travail utilise correctement MPI4Py ou Dask.
  • La reproductibilité de l’environnement est souvent la partie la plus difficile du travail de HPC. Les environnements Conda, les fichiers de verrouillage et les scripts de travaux aident à le résoudre.
  • Le flux de travail comporte trois étapes : prototyper localement, paralléliser avec MPI4PY ou DASK et soumettre des travaux via SLURM, PBS ou un autre planificateur.
  • Ne réécrivez pas toute la base de code pour le cluster. Gardez la logique de simulation stable et enveloppez-la avec un environnement reproductible et des couches d’exécution parallèles.

Vous écrivez un script de simulation sur votre ordinateur portable. Vous le testez avec de petits jeux de données. Il fonctionne en quelques minutes. Ensuite, vous avez besoin de milliers de cœurs et de centaines de gigaoctets de RAM pour l’exécuter à grande échelle.

Réécrivez-vous toute la base de code ?

Non. Le même code Python qui s’exécute sur votre ordinateur portable peut s’exécuter sur un superordinateur avec des modifications minimales. La clé n’est pas de réécrire. C’est enveloppant.

Vous avez besoin de trois éléments : un environnement reproductible, un modèle d’exécution parallèle et un script de soumission de travaux pour le planificateur de cluster.

Ce guide parcourt le flux de travail complet du prototype local au déploiement HPC de production sans modifier votre logique de simulation de base.

Le flux de travail HPC Python en trois étapes

La plupart des projets HPC Python suivent le même modèle en trois étapes, quel que soit le domaine scientifique.

Stage 1 — Local Prototyping:
  Laptop → Jupyter or IDE → NumPy / Pandas → serial execution

Stage 2 — Parallelization:
  Same code → mpi4py or Dask → parallel execution

Stage 3 — Cluster Deployment:
  Python script → Slurm / PBS job submission → distributed compute nodes

L’objectif est de parcourir ces étapes progressivement. Commencez avec une version locale fonctionnelle. Ajoutez une exécution parallèle uniquement une fois que la logique est correcte. Soumettre au cluster uniquement après des petits tests parallèles.

Étape 1 : Mise en place d’un environnement reproductible

Avant d’écrire du code parallèle, vous avez besoin d’un environnement de développement fiable. C’est là que de nombreux flux de travail scientifiques échouent en premier.

Pourquoi la reproductibilité est difficile pour les clusters HPC

Les clusters HPC sont des systèmes partagés. De nombreux groupes de recherche utilisent la même infrastructure et les packages de systèmes peuvent changer au fil du temps. Une simulation qui fonctionne aujourd’hui peut interrompre six mois plus tard après les mises à jour du module, les modifications du compilateur ou les modifications de la version du package.

La solution est la gestion de l’environnement. Conda est couramment utilisé car il isole les packages du système Python et peut gérer les dépendances compilées ainsi que les packages Python.

# Create a reproducible environment
conda create -n my-sim python=3.11 numpy mpi4py dask

# Lock dependencies for reproducibility
conda env export --no-builds --name my-sim > environment.yml

# Recreate the environment on another machine
conda env create -n my-sim -f environment.yml

Un seul fichier d’environnement donne aux collaborateurs et aux futurs utilisateurs un moyen clair de recréer la pile logicielle. Pour une reproductibilité plus stricte, utilisez des fichiers de verrouillage qui épinglent chaque version de package et dépendance de génération.

Conda vs Venv : pourquoi Conda gagne souvent sur HPC

Python intégré venv fonctionne bien pour les packages Python Pure. Les flux de travail HPC dépendent souvent de bibliothèques scientifiques compilées, d’exécutions MPI, de HDF5, de BLA, de CUDA et d’autres composants au niveau du système.

Fonctionnalité veuve peste
Forfaits Pure Python Oui Oui
Dépendances au niveau du système telles que MPI et HDF5 Non Oui
Cohérence multiplateforme Limité Fort
Forfaits d’accélération GPU tels que Cupy ou PyTorch CUDA Configuration manuelle Configuration plus automatisée

Conda donne une reproductibilité dans une plus grande partie de la pile scientifique. Cela compte lorsqu’il travaille avec MPI4Py car l’environnement d’exécution MPI et Python doit être compatible à la fois sur la machine locale et sur le cluster.

Étape 2 : Parallèlement à votre code Python

Une fois que la version locale fonctionne, la prochaine étape est la parallélisation. Dans les flux de travail HPC Python, deux outils sont particulièrement courants :

  1. MPI4py pour un contrôle distribué à grain fin sur de nombreux nœuds.
  2. DASK pour un parallélisme de niveau supérieur avec moins de changements de code.

Pourquoi vous avez besoin de MPI4PY

Le verrou d’interpréteur global de Python limite la véritable exécution de Python multithread dans un processus. Le module multiprocessing peut paralléliser sur une seule machine, mais il ne s’adapte pas naturellement à plusieurs nœuds de calcul.

MPI4PY fournit des liaisons Python pour MPI, l’interface de transmission de messages. MPI est une API standard pour le calcul parallèle distribué et est largement utilisé sur les systèmes HPC.

Avec MPI4Py, chaque processus exécute le même script mais reçoit un rang unique. Ce classement contrôle la partie de la charge de travail que chaque processus gère.

Modèle MPI4PY de base : bonjour le monde

from mpi4py import MPI

comm = MPI.COMM_WORLD

rank = comm.Get_rank()
size = comm.Get_size()

print(f"Hello from process {rank} out of {size} processes")

Lancez le script avec MPI :

mpiexec -n 16 python my_script.py

Cela démarre 16 processus Python indépendants. Chaque processus exécute le même script, mais chacun a un rank différent.

Le modèle de données : processus indépendants

Les processus MPI ne partagent pas de mémoire par défaut. Chaque processus a son propre espace mémoire. Pour échanger des informations, les processus doivent explicitement envoyer, recevoir, diffuser, collecter ou réduire les données.

import numpy as np
from mpi4py import MPI

comm = MPI.COMM_WORLD
rank = comm.Get_rank()

# Rank 0 creates the initial data.
if rank == 0:
    data = np.arange(10, dtype="i")
else:
    data = np.empty(10, dtype="i")

# Broadcast data from rank 0 to all ranks.
comm.Bcast(data, root=0)

print(f"Rank {rank} received data: {data}")

Ce modèle de communication explicite est l’une des raisons pour lesquelles le MPI évolue bien. Chaque classement possède sa mémoire locale et la communication ne se produit que lorsque vous en faites la demande.

Distribution de la charge de travail : le modèle de base

La plupart des flux de travail MPI scientifiques suivent le même schéma : divisez le problème, calculez localement et réduisez les résultats.

from mpi4py import MPI
import numpy as np

comm = MPI.COMM_WORLD

size = comm.Get_size()
rank = comm.Get_rank()

# Total problem size
N = 10_000_000

# Calculate workload per rank
workloads = [N // size for _ in range(size)]

for i in range(N % size):
    workloads[i] += 1

my_start = sum(workloads[:rank])
my_end = my_start + workloads[rank]

# Each rank works on its own slice
my_data = np.random.rand(my_end - my_start)

# Local computation
local_result = np.sum(np.sin(my_data))

# Sum results across all ranks
send_buffer = np.array([local_result])
receive_buffer = np.zeros(1)

comm.Reduce(send_buffer, receive_buffer, op=MPI.SUM, root=0)

if rank == 0:
    print(f"Total computed across {size} ranks: {receive_buffer[0]}")

Ce modèle s’applique à de nombreuses tâches de simulation :

  • Simulations de Monte Carlo, où chaque classement traite des échantillons indépendants.
  • Balayages de paramètres, chaque classement testant différents jeux de paramètres.
  • Décomposition du domaine, où chaque classement possède une partie de la grille de simulation.
  • Dynamique moléculaire, où les rangs calculent les forces pour différents sous-ensembles de particules.

La boîte à outils de communication collective

MPI4PY fournit des opérations collectives qui sont généralement plus faciles et plus efficaces que la messagerie point à point personnalisée.

Opération ce qu’il fait Quand utiliser
Bcast Un processus envoie les mêmes données à tous les rangs Partage des conditions initiales, des constantes ou des valeurs de configuration
Scatter Distribue des morceaux d’un tableau à travers les rangs Décomposition du domaine ou partitionnement de la charge de travail
Gather Collecte des données de tous les rangs Collecte des sorties partielles
Reduce Agrége des valeurs telles que somme, max ou min Sommer des énergies ou collecter des métriques globales
Allreduce Agrége les valeurs et donne le résultat à chaque rang Synchroniser l’état global dans tous les processus

Utilisez la communication collective lorsque cela est possible. Il est généralement plus simple et plus efficace que de coordonner manuellement de nombreux envois et réceptions.

Quand utiliser DASK à la place

Dask se situe entre Python en série et MPI. Il est utile lorsque la charge de travail est parallèlement embarrassante, comme l’exécution de la même simulation avec de nombreux jeux de paramètres ou conditions initiales.

from dask.distributed import Client
from dask import delayed

client = Client(n_workers=16, threads_per_worker=4)

def my_simulation(params):
    # Your simulation logic here
    return params["a"] * params["b"]

parameter_list = [
    {"a": i, "b": 2.0}
    for i in range(100)
]

tasks = [
    delayed(my_simulation)(params)
    for params in parameter_list
]

final_results = client.compute(tasks, sync=True)

print(final_results[:5])

L’avantage de Dask est qu’il vous permet de paralléliser de nombreux flux de travail sans redéfinir la simulation autour des rangs et des communicateurs.

MPI4py vs DASK : lequel choisir ?

Critère Mpi4py ténèbres
Courbe d’apprentissage Plus raide parce que vous devez comprendre les rangs et les communicateurs Plus doux car il utilise des modèles de tâches Python familiers
Contrôle à grain fin Excellent car chaque communication est explicite Limité par l’abstraction du graphe de tâches
Efficacité mémoire Élevé car chaque classement stocke sa propre tranche Modéré parce que les travailleurs ajoutent des frais généraux
Convient à Solveurs PDE à grande échelle et simulations décomposées par le domaine Monte Carlo, balayages de paramètres et analyse des données
Meilleure échelle Parallélisme dense sur de nombreux nœuds Échelle faible sur le nombre de travailleurs modérés

Commencez par Dask pour le prototypage et des charges de travail parallèles embarrassantes. Passez à MPI4PY lorsque vous avez besoin d’un contrôle étroit de la communication, de la mémoire et de la mise à l’échelle sur de nombreux nœuds.

Étape 3 : Soumission au cluster

Après les tests et la parallélisation locaux, la prochaine étape est le déploiement de cluster. La plupart des clusters utilisent un planificateur de travaux. Le slurm est l’un des plus courants.

Le script de travail Slurm

Un script de tâche Slurm indique au planificateur quelles ressources vous avez besoin et comment exécuter le code.

#!/bin/bash
#SBATCH --job-name=my-sim
#SBATCH --nodes=32
#SBATCH --tasks-per-node=4
#SBATCH --cpus-per-task=1
#SBATCH --mem=8GB
#SBATCH --time=04:00:00
#SBATCH --output=sim_output.%j

# Load required modules
module load Python/3.11

# Activate conda environment
eval "$(conda shell.bash hook)"
conda activate my-sim

# Run the MPI job
srun python my_simulation.py

Les paramètres importants comprennent :

  • --nodes : nombre de nœuds de calcul.
  • --tasks-per-node : nombre de tâches MPI par nœud.
  • --cpus-per-task : cœurs de CPU affectés à chaque tâche.
  • --time : limite d’horloge murale. Les emplois sont généralement arrêtés s’ils le dépassent.
  • --output : modèle de fichier de sortie. %j Insère l’ID du travail.

Exécution sur des clusters sans Slurm

Tous les clusters n’utilisent pas Slurm. Les alternatives courantes comprennent :

  • PBS ou couple, généralement en utilisant qsub.
  • LSF, en utilisant généralement bsub.
  • Cobalt ou autres planificateurs spécifiques au site.

La syntaxe de la soumission des travaux change, mais le code Python ne change généralement pas. MPI4PY et DASK sont pour la plupart agnostiques de planification une fois lancés correctement.

Mode interactif vs batch

Pour le débogage, les exécutions interactives sont utiles :

srun --ntasks=4 --pty --time=02:00:00 python my_simulation.py

Pour la production, utilisez la soumission par lot :

sbatch my_job_script.sh

Les sessions interactives sont bonnes pour les tests courts. Les travaux par lots sont meilleurs pour les longues périodes, les charges de travail du jour au lendemain et les simulations de production.

Le flux de travail complet : de l’ordinateur portable à la production

La transition du prototype local au déploiement du cluster peut être progressive. La logique de simulation de base doit rester stable pendant que l’enveloppe d’exécution change.

Étape 1 : Développer et tester localement

# my_simulation.py
import numpy as np

def simulate(initial_condition, params):
    # Core simulation logic
    result = np.sin(initial_condition) * params["factor"]
    return np.sum(result)

# Local test
data = np.random.rand(1000)
result = simulate(data, {"factor": 1.5})

print(f"Local result: {result}")

Étape 2 : Ajout de la parallélisation

# my_simulation_mpi.py
import numpy as np
from mpi4py import MPI

comm = MPI.COMM_WORLD

rank = comm.Get_rank()
size = comm.Get_size()

def simulate(initial_condition, params):
    # Core simulation logic stays the same
    result = np.sin(initial_condition) * params["factor"]
    return np.sum(result)

N = 10_000_000

workloads = [N // size for _ in range(size)]

for i in range(N % size):
    workloads[i] += 1

my_start = sum(workloads[:rank])
my_end = my_start + workloads[rank]

my_data = np.random.rand(my_end - my_start)

local_result = simulate(my_data, {"factor": 1.5})

send_buffer = np.array([local_result])
receive_buffer = np.zeros(1)

comm.Reduce(send_buffer, receive_buffer, op=MPI.SUM, root=0)

if rank == 0:
    print(f"Parallel result across {size} ranks: {receive_buffer[0]}")

Étape 3 : Soumettre au cluster

Enregistrez un script de travail tel que job_script.sh :

#!/bin/bash
#SBATCH --job-name=parallel-sim
#SBATCH --nodes=64
#SBATCH --tasks-per-node=8
#SBATCH --time=12:00:00

module load Python/3.11

eval "$(conda shell.bash hook)"
conda activate my-sim

srun python my_simulation_mpi.py

Soumettez-le :

sbatch job_script.sh

La logique de simulation reste la même. Le wrapper modifie la façon dont la charge de travail est divisée et lancée.

Des pièges courants et comment les éviter

Pitfall 1 : Diffuser trop de données

La diffusion de grands tableaux de rang 0 à chaque rang peut devenir un goulot d’étranglement de communication. Pour les simulations volumineuses, évitez d’envoyer plus de données que ce dont chaque rang a besoin.

Les meilleures options incluent :

  • Utilisez Scatter au lieu de Bcast lorsque chaque rang n’a besoin que d’une tranche.
  • Utilisez la décomposition du domaine afin que chaque rang possède une région de la simulation.
  • Pour Monte Carlo, laissez chaque rang générer ses propres échantillons aléatoires au lieu de les diffuser.

Pitfall 2 : surallocation des ressources du cluster

Demander plus de nœuds que votre code peut utiliser des pertes de temps d’allocation et peut augmenter le délai de file d’attente.

Commencez par un petit travail de test, tel que 4 à 8 nœuds. Mesurez l’accélération. Échelle uniquement après que le code montre une efficacité parallèle utile.

Pitfall 3 : Oublier le Gil en code mixte

Les bibliothèques NumPy, Scipy et Compiled C ou Fortran libèrent souvent le verrou d’interpréteur global Python lors d’opérations numériques lourdes. Les threads Pure Python ne procurent généralement pas le même avantage.

Si votre flux de travail mélange les threads Python avec les bibliothèques numériques, testez la mise à l’échelle au lieu de supposer que plus de threads vous aideront.

Pitfall 4 : inadéquation de l’environnement sur le cluster

Votre ordinateur portable peut utiliser Python 3.11 tandis que le module de cluster fournit Python 3.9. Les bibliothèques MPI4Py, MPI et les dépendances compilées peuvent se casser lorsque les versions ne correspondent pas.

Utilisez un environnement Conda et documentez la configuration exacte. Recréez et testez l’environnement sur le cluster avant d’exécuter des travaux volumineux.

Guide de décision : quelle stratégie de parallélisation ?

Situation Approche recommandée
Petit ensemble de données et tests de logique locale Python série avec NumPy
Machine unique avec plusieurs cœurs Python multiprocessing ou DASK LocalCluster
8 à 64 travailleurs et tâches parallèles embarrassantes DASK avec un lanceur de cluster
64 à 1 000 nœuds et PDE décomposés dans le domaine MPI4Py avec une distribution de charge de travail explicite
Simulation de production sur plus de 1 000 noeuds MPI4Py, scripts de travaux de planificateur et environnement Conda verrouillé
Code de recherche sans reproductibilité garantie DASK, fichier d’environnement Conda et script de travail comme point de départ

Votre outil de parallélisation doit correspondre à l’échelle et à la structure du problème. Commencez plus petit que ce dont vous pensez avoir besoin, puis augmentez après la mesure des performances.

Guides connexes

Pour les sujets connexes dans les flux de travail de simulation scientifique :

Résumé et étapes suivantes

Les flux de travail HPC Python concernent l’emballage, et non la réécriture. La même logique de simulation qui s’exécute sur un ordinateur portable peut s’exécuter sur de nombreux cœurs lorsque vous créez la bonne structure d’exécution.

Le workflow est :

  1. Verrouillez l’environnement avec Conda afin que le code s’exécute de manière cohérente entre les systèmes.
  2. Ajoutez du parallélisme avec MPI4PY pour un contrôle distribué à grain fin ou un DASK pour un parallélisme de tâches de haut niveau.
  3. Soumettez des travaux via Slurm, PBS, LSF ou le planificateur utilisé par votre cluster.

La partie la plus difficile est souvent l’environnement, pas le code. Investissez rapidement dans une configuration reproductible et les simulations parallèles deviennent plus faciles à mettre à l’échelle de l’ordinateur portable au superordinateur.

Prochaines étapes

  1. Vérifier la simulation actuelle. S’il ne fonctionne que sur un ordinateur portable, commencez par créer un environnement Conda et tester une petite exécution multicœur.
  2. Mesurez avant de mettre à l’échelle. Exécutez un petit travail de cluster, mesurez l’accélération, puis passez à la taille de la production.
  3. Documentez la configuration. Enregistrez le fichier d’environnement, le script de travail, la commande de lancement et la sortie attendue.

L’établissement d’un flux de travail HPC Python nécessite de comprendre les packages scientifiques Python, la programmation parallèle et l’infrastructure de cluster. Si vous avez besoin d’aide pour la parallélisation MPI, la planification des tâches ou la reproductibilité de l’environnement, notre équipe peut prendre en charge des flux de travail scientifiques évolutifs de Python, notamment des simulations basées sur Fipy et des moteurs de Monte Carlo personnalisés.