TRAVAILLEUR_NCCL = r'''
import json
import os
import time
import torch
import torch.nn as nn
import torch.distributed as dist
from functools import partial
from torch.distributed.fsdp import FullyShardedDataParallel as FSDP
from torch.distributed.fsdp import ShardingStrategy
from torch.distributed.fsdp.wrap import size_based_auto_wrap_policy
RANG = int(os.environ["RANK"])
MONDE = int(os.environ["WORLD_SIZE"])
STRATEGIE = os.environ["STRATEGIE"]
RALENTI = float(os.environ.get("RALENTI", "0")) # secondes de pause par pas, rang 0 seulement
COUCHES, LARGEUR, LOT, PAS = 8, 4096, 64, 6
def octets(tenseurs):
return sum(t.numel() * t.element_size() for t in tenseurs if torch.is_tensor(t))
def modele(graine=0):
torch.manual_seed(graine)
return nn.Sequential(*[nn.Linear(LARGEUR, LARGEUR) for _ in range(COUCHES)])
def main():
torch.cuda.set_device(RANG)
periph = torch.device("cuda", RANG)
dist.init_process_group("nccl", init_method=os.environ["DIST_INIT"],
rank=RANG, world_size=MONDE)
base = modele().to(periph)
if STRATEGIE == "ddp":
m = nn.parallel.DistributedDataParallel(base, device_ids=[RANG])
opt = torch.optim.AdamW(m.parameters(), lr=1e-3)
elif STRATEGIE == "fsdp":
m = FSDP(base, sharding_strategy=ShardingStrategy.FULL_SHARD,
auto_wrap_policy=partial(size_based_auto_wrap_policy,
min_num_params=10**6),
device_id=periph)
opt = torch.optim.AdamW(m.parameters(), lr=1e-3)
else:
raise SystemExit("strategie inconnue : " + STRATEGIE)
torch.cuda.reset_peak_memory_stats()
temps_par_pas = []
for pas in range(PAS):
debut = time.perf_counter()
x = torch.randn(LOT, LARGEUR, device=periph)
opt.zero_grad(set_to_none=True)
y = m(x)
y.sum().backward()
opt.step()
torch.cuda.synchronize()
fin = time.perf_counter()
temps_par_pas.append(fin - debut)
if RANG == 0 and RALENTI:
time.sleep(RALENTI) # hors fenetre mesuree : l'attente paiera chez l'autre
poids = octets(list(m.parameters()))
optim = octets([v for etat in opt.state.values() for v in etat.values()])
sortie = {"rang": RANG, "gpu": torch.cuda.get_device_name(RANG),
"strategie": STRATEGIE, "ralenti": RALENTI,
"poids": poids, "optimiseur": optim,
"pic_cuda_mo": round(torch.cuda.max_memory_allocated() / 2**20, 1),
"pas_ms": [round(t * 1000, 1) for t in temps_par_pas]}
print("RESULTAT " + json.dumps(sortie), flush=True)
dist.destroy_process_group()
if __name__ == "__main__":
main()
'''3.13b — Deux vraies cartes : DDP et FSDP sous nccl, et le rang qui fait attendre l’autre
Navigation : << 3.13 — Découper le modèle : DDP, ZeRO, FSDP · Feuille de route de la série
Le carnet 3.13 a découpé un modèle huit fois — DDP, ZeRO-1/2/3, pipelines — mais tout sur CPU, sous gloo, avec des octets comptés à la main. Il reste le geste que CPU ne peut pas faire : voir nccl à l’œuvre, mesurer la mémoire réelle des cartes (max_memory_allocated, pas un calcul), chronométrer le pas — et lire la signature du défaut le plus courant du multi-GPU : « ma carte attend l’autre ».
Le dispositif est volontairement le même que 3.13 — le même modèle, le même partage du travail — pour qu’une seule chose change entre les deux carnets : le backend et le matériel.
Où ce carnet s’exécute.
nccln’existe pas sous Windows : ce carnet tourne dans un conteneur Linux (image PyTorch officielle) avec les deux cartes passées au conteneur. Les sorties committées viennent de cette exécution ; la machine et les cartes exactes sont nommées dans la première cellule de code.
Le dispositif : un modèle, deux rangs, deux cartes
Même modèle qu’en 3.13 — un empilement de couches linéaires — mais dimensionné pour que la carte le sente : 8 couches de 4096 (environ 134 millions de paramètres, un demi-milliard d’octets en simple précision). Le monde est de deux processus, un rang par carte. Les deux cartes de la machine d’exécution sont hétérogènes — c’est un avantage pédagogique : l’asymétrie naturelle qu’on mesure au pas 1 est exactement celle qu’on cherchera à amplifier exprès au diagnostic du rang lent.
Le lanceur
Même idiome qu’en 3.13 : le travailleur vit dans une chaîne, le lanceur écrit un fichier, démarre un processus par rang (RANK dans l’environnement), et ramasse les lignes RESULTAT. Seule différence : chaque rang s’installe sur sa carte (torch.cuda.set_device) au démarrage.
import json
import os
import pathlib
import subprocess
import sys
import tempfile
import torch
def carte_de_travail():
"""Verifie ce que le conteneur voit avant toute Experience : deux cartes, nccl pret."""
assert torch.cuda.device_count() >= 2, (
"il faut deux cartes visibles (conteneur lance avec --gpus all)")
assert torch.distributed.is_nccl_available(), "nccl indisponible dans cette image"
for r in range(torch.cuda.device_count()):
print("carte %d : %s" % (r, torch.cuda.get_device_name(r)))
def lance(monde, travailleur, strategie, ralenti=0.0, delai=300):
"""Lance `monde` processus du travailleur et ramasse leurs lignes RESULTAT."""
dossier = pathlib.Path(tempfile.mkdtemp(prefix="nccl_"))
chemin = dossier / "travailleur.py"
chemin.write_text(travailleur, encoding="utf-8")
store = pathlib.Path(tempfile.mkdtemp(prefix="nccl_store_"))
init = "file:///" + str(store / "rv").replace("\\", "/")
env = dict(os.environ, DIST_INIT=init, WORLD_SIZE=str(monde),
STRATEGIE=strategie, RALENTI=str(ralenti))
procs = [subprocess.Popen([sys.executable, str(chemin)],
env=dict(env, RANK=str(r)),
stdout=subprocess.PIPE, stderr=subprocess.PIPE,
text=True, encoding="utf-8") for r in range(monde)]
resultats, pannes = [], []
for rang, proc in enumerate(procs):
sortie, erreur = proc.communicate(timeout=delai)
if proc.returncode != 0:
pannes.append("rang %d : code %d\n%s" % (rang, proc.returncode, erreur.strip()[-600:]))
for ligne in sortie.splitlines():
if ligne.startswith("RESULTAT "):
resultats.append(json.loads(ligne[len("RESULTAT "):]))
if pannes:
raise RuntimeError("echec de %s a monde %d :\n%s" % (strategie, monde, "\n".join(pannes)))
return sorted(resultats, key=lambda r: r["rang"])
carte_de_travail()carte 0 : NVIDIA GeForce RTX 3090
carte 1 : NVIDIA GeForce RTX 3080 Ti Laptop GPU
Premier dépouillement : ce que chaque stratégie occupe sur les cartes
Deux runs : DDP (chaque rang porte tout) puis FSDP FULL_SHARD (chaque rang porte sa part). Pour chaque run : les octets détenus par chaque rang — le même comptage qu’en 3.13 — plus deux mesures que CPU ne pouvait pas donner : le pic de mémoire CUDA (max_memory_allocated, tout ce que la carte a réellement vécu, activations comprises) et le pas médian en millisecondes.
COUCHES, LARGEUR = 8, 4096
P = COUCHES * (LARGEUR * LARGEUR + LARGEUR) # parametres du modele, comme en 3.13
releves = {}
for strategie in ("ddp", "fsdp"):
rangs = lance(2, TRAVAILLEUR_NCCL, strategie)
releves[strategie] = rangs
for r in rangs:
print("monde=2 %-5s rang%d %-28s poids=%3.2f P pic=%6.0f Mo pas median=%6.1f ms"
% (strategie, r["rang"], r["gpu"][:28],
r["poids"] / P, r["pic_cuda_mo"],
sorted(r["pas_ms"])[len(r["pas_ms"]) // 2]))monde=2 ddp rang0 NVIDIA GeForce RTX 3090 poids=4.00 P pic= 3091 Mo pas median= 401.7 ms
monde=2 ddp rang1 NVIDIA GeForce RTX 3080 Ti L poids=4.00 P pic= 3091 Mo pas median= 402.5 ms
monde=2 fsdp rang0 NVIDIA GeForce RTX 3090 poids=2.00 P pic= 1299 Mo pas median= 654.2 ms
monde=2 fsdp rang1 NVIDIA GeForce RTX 3080 Ti L poids=2.00 P pic= 1299 Mo pas median= 655.0 ms
Lecture : les octets confirment 3.13, les cartes ajoutent deux colonnes
Les rapports d’octets sont exactement ceux de 3.13 — DDP détient \(4P\) par rang (poids \(P\) + gradients \(P\) + optimiseur \(2P\)), FSDP FULL_SHARD retombe à \(2P\) (chaque part divisée par le monde). Le backend ne change rien au découpage : nccl transport, ZeRO découpe — deux problèmes indépendants, et ce run sur vraies cartes le vérifie au lieu de l’affirmer.
Les deux colonnes nouvelles : le pic CUDA (plus haut que le seul état détenu — les activations et les tampons d’all-gather vivent aussi sur la carte) et le pas — où l’asymétrie native des deux cartes se voit déjà, avant même d’organiser quoi que ce soit.
Le diagnostic : « ma carte attend l’autre »
Troisième run, DDP, avec une pause de 50 ms après chaque pas sur le rang 0 seulement — et hors de la fenêtre mesurée. Personne ne ralentit son propre travail : le rang 0 dort entre deux pas. Mais le all-reduce du pas suivant ne se termine que lorsque tout le monde y arrive — donc l’attente doit apparaître… chez le rang sain. C’est toute la signature du défaut : celui qui se plaindra d’une carte lente n’est pas celui qui la cause.
rangs_lents = lance(2, TRAVAILLEUR_NCCL, "ddp", ralenti=0.05)
rapide, lent = rangs_lents[1], rangs_lents[0] # rang 1 sans pause, rang 0 avec
print("rang 0 (%s) : pause de 50 ms apres chaque pas, hors mesure" % lent["gpu"][:24])
print("rang 1 (%s) : rien a signaler" % rapide["gpu"][:24])
print()
print("%4s %10s %10s" % ("pas", "rang0 ms", "rang1 ms"))
for k, (t0, t1) in enumerate(zip(lent["pas_ms"], rapide["pas_ms"])):
barre = " " + "#" * int(t1 / 10)
print("%4d %10.1f %10.1f%s" % (k + 1, t0, t1, barre))rang 0 (NVIDIA GeForce RTX 3090) : pause de 50 ms apres chaque pas, hors mesure
rang 1 (NVIDIA GeForce RTX 3080 ) : rien a signaler
pas rang0 ms rang1 ms
1 709.9 718.1 #######################################################################
2 395.2 446.5 ############################################
3 393.6 443.7 ############################################
4 394.2 444.3 ############################################
5 394.3 444.4 ############################################
6 392.9 443.1 ############################################
Lecture : l’attente se lit chez celui qui ne dort pas
Le rang 0 dort mais ses pas mesurés restent courts — la pause est hors fenêtre et son all-reduce trouve le rang 1 déjà arrivé. Le rang 1, lui, n’a changé rien — et ses pas ont grossi d’à-peu-près la pause de l’autre : son all-reduce attend le dormant. En pratique, multiplier par le nombre de pas : la carte saine d’un couple déséquilibré perd le temps exact que l’autre gaspille — et le sommet de la courbe de la carte « lente » en production n’est souvent pas le coupable, c’est la victime.
Exercices
Trois exercices, dans l’esprit de la série : prédire, puis vérifier contre le relevé committé ci-dessus.
# Exercice 1 -- Le pic de FSDP.
# Predire le pic CUDA (en Mo) du run FSDP, a partir de P et du releve DDP :
# l'etat detenu passe de 4P a 2P par rang, mais l'all-gather regroupe chaque
# couche a la volee. Ecrire l'estimation, puis la comparer au releve ci-dessus.
# Indice : le pic FSDP n'est pas la moitie du pic DDP -- les tampons temporaires
# de regroupement pesent le poids d'une couche entiere, pas d'un fragment.
# TODO etudiant
print("Exercice a completer : estimation du pic FSDP, comparaison au releve.")Exercice a completer : estimation du pic FSDP, comparaison au releve.
# Exercice 2 -- Inverser la pause.
# Si RALENTI s'appliquait au rang 1 au lieu du rang 0 (changer la condition dans
# le travailleur, relancer), qu'afficheraient les deux colonnes du tableau ?
# Ecrire la prediction AVANT de relancer, puis verifier.
# Etape 1 : ecrire la prediction ici.
# Etape 2 : modifier la condition, relancer la cellule du diagnostic, comparer.
# TODO etudiant
print("Exercice a completer : prediction de la colonne inversee.")Exercice a completer : prediction de la colonne inversee.
# Exercice 3 -- Le cout du ralenti.
# Avec une pause de 50 ms sur 6 pas, estimer le temps total PERDU par le rang 1
# (celui qui attend), puis le temps total perdu par l'entrainement entier si le
# monde passe de 2 a 8 rangs avec un seul lent. Verifier contre le tableau.
# TODO etudiant
print("Exercice a completer : cout du ralenti, a 2 puis a 8 rangs.")Exercice a completer : cout du ralenti, a 2 puis a 8 rangs.
Suite
- 3.13 — Découper le modèle : le socle — les quatre stratégies, les octets par rang, le pipeline, sur CPU sous
gloo; - les pipelines ML-Training-Pipeline (QuantConnect) consomment ce geste à l’échelle des vrais entraînements multi-cartes ;
- la série GenAI vit l’autre face du multi-GPU : l’inférence découpée (offload), où les mêmes cartes portent des poids au lieu de gradients.