Parallélisme et concurrence
Mis à jour le 28 juillet 2026
Attendre à plusieurs, ce n’est pas calculer à plusieurs
Un appel à l’API OpenAI est une opération I/O-bound : votre programme ne calcule rien, il attend une réponse réseau. Cette nature dicte le bon outil. Le parallélisme, au sens du multiprocessing, sert à répartir du calcul sur plusieurs cœurs ; il coûte un processus complet par unité de travail, et faire tourner cent processus qui attendent tous le réseau est un gaspillage pur. La concurrence, telle que la fournit asyncio, permet à un seul processus de tenir des centaines d’attentes simultanées pour un coût dérisoire. C’est elle qui maximise le débit de vos appels.
La distinction est plus qu’académique : une équipe qui parallélise ses appels LLM avec un pool de processus consommera plusieurs gigaoctets de mémoire pour un gain nul par rapport à quelques lignes d’asyncio.
Lancer des appels en concurrence
Le SDK expose un client asynchrone, AsyncOpenAI, dont les appels s’attendent avec await. Rassembler plusieurs coroutines dans un asyncio.gather les fait progresser toutes en même temps, la durée totale se rapprochant de celle du plus lent appel plutôt que de leur somme.
import asyncio
import openai
client = openai.AsyncOpenAI()
async def appel_unique(prompt: str, modele: str = "gpt-5.6-terra") -> str:
"""Un appel API asynchrone."""
response = await client.responses.create(
model=modele,
input=prompt,
)
return response.output_text
async def appels_concurrents(prompts: list[str]) -> list[str]:
"""Lance tous les appels en concurrence."""
taches = [appel_unique(prompt) for prompt in prompts]
resultats = await asyncio.gather(*taches, return_exceptions=True)
for i, resultat in enumerate(resultats):
if isinstance(resultat, Exception):
print(f"Erreur sur le prompt {i}: {resultat}")
return [r for r in resultats if not isinstance(r, Exception)]
Le return_exceptions=True est ce qui distingue un prototype d’un traitement de production. Sans lui, la première exception interrompt le gather et vous perdez les quatre-vingt-dix-neuf résultats déjà obtenus. Avec lui, les erreurs deviennent des valeurs que vous inspectez, journalisez et retentez à votre rythme.
Cette version brute a toutefois un défaut immédiat : lancer mille appels d’un coup vous fera franchir les limites de débit dans la seconde. Un sémaphore résout le problème en n’autorisant qu’un nombre fixe d’appels simultanés, les autres attendant sagement leur tour. Vingt est un point de départ raisonnable, à ajuster selon les quotas réels de votre organisation.
async def appels_avec_limite(
prompts: list[str],
max_concurrence: int = 20,
) -> list[str]:
"""Limite le nombre d'appels simultanés."""
semaphore = asyncio.Semaphore(max_concurrence)
async def appel_limite(prompt: str) -> str:
async with semaphore:
return await appel_unique(prompt)
taches = [appel_limite(prompt) for prompt in prompts]
return await asyncio.gather(*taches, return_exceptions=True)
Le pattern map-reduce appliqué au texte
Un document de deux cents pages ne se résume pas en un appel, et le découper séquentiellement prendrait des minutes. Le pattern map-reduce s’y applique naturellement : on découpe en fragments, on résume chaque fragment en concurrence, puis on demande au modèle une synthèse des synthèses. La phase map profite pleinement de la concurrence, la phase reduce est un appel unique sur un texte devenu court.
async def map_reduce_resume(
document: str,
taille_chunk: int = 3000,
) -> str:
"""Résume un long document par map-reduce parallèle."""
# Phase MAP : découper et résumer en parallèle
chunks = [
document[i:i + taille_chunk]
for i in range(0, len(document), taille_chunk)
]
prompts_map = [
f"Résumez ce passage en 3 points clés :\n\n{chunk}"
for chunk in chunks
]
resumes_partiels = await appels_avec_limite(prompts_map, max_concurrence=10)
# Phase REDUCE : combiner les résumés
resumes_texte = "\n\n---\n\n".join(
r for r in resumes_partiels if isinstance(r, str)
)
resume_final = await appel_unique(
f"Synthétisez ces résumés partiels en un résumé cohérent "
f"de 5 points clés :\n\n{resumes_texte}"
)
return resume_final
Le filtrage if isinstance(r, str) est une décision produit autant qu’une précaution technique : il signifie que la synthèse finale sera produite même si trois fragments ont échoué, en silence. Selon le contexte, c’est exactement ce que vous voulez, ou c’est inacceptable. Décidez explicitement, et journalisez le nombre de fragments perdus dans tous les cas.
Quand le multiprocessing redevient pertinent
Il reste des situations, rares, où le calcul reprend ses droits : un post-traitement lourd sur les réponses, extraction d’entités, scoring, parsing structuré. Ce travail-là est CPU-bound et ne gagne rien à l’asyncio. La combinaison élégante consiste à enchaîner les deux régimes dans un même pipeline : les appels réseau en concurrence, puis le calcul réparti sur un pool de processus via run_in_executor.
import asyncio
from concurrent.futures import ProcessPoolExecutor
def post_traitement_lourd(texte: str) -> dict:
"""Traitement CPU-bound sur le résultat (NLP, parsing, etc.)."""
# Exemple : extraction d'entités, scoring, etc.
mots = texte.split()
return {
"nb_mots": len(mots),
"mots_uniques": len(set(mots)),
}
async def pipeline_complet(prompts: list[str]) -> list[dict]:
"""Pipeline : appels LLM async + post-traitement CPU parallèle."""
# Étape 1 : appels LLM en concurrence
resultats_llm = await appels_avec_limite(prompts)
# Étape 2 : post-traitement en parallèle (CPU)
loop = asyncio.get_event_loop()
with ProcessPoolExecutor(max_workers=4) as pool:
taches = [
loop.run_in_executor(pool, post_traitement_lourd, r)
for r in resultats_llm
if isinstance(r, str)
]
resultats_finaux = await asyncio.gather(*taches)
return resultats_finaux
Vérifiez que le calcul justifie vraiment le pool avant de l’introduire : sérialiser les données vers un autre processus a un coût, et sur des traitements de quelques millisecondes ce coût dépasse le gain.
Vérifier le débit réel
Le réglage de max_concurrence ne se devine pas. Trop bas, vous laissez du débit sur la table ; trop haut, vous déclenchez des limitations qui vous coûtent plus de temps qu’elles ne vous en font gagner. Mesurez sur votre charge réelle, en faisant varier la valeur, et regardez à la fois le débit et le nombre de succès — un pipeline très rapide qui échoue une fois sur cinq n’est pas rapide.
import time
async def benchmark_debit(nb_requetes: int = 100):
"""Mesure le débit réel de votre pipeline."""
prompts = [
f"Donnez un fait intéressant numéro {i}."
for i in range(nb_requetes)
]
debut = time.perf_counter()
resultats = await appels_avec_limite(prompts, max_concurrence=20)
duree = time.perf_counter() - debut
succes = sum(1 for r in resultats if isinstance(r, str))
print(f"Requêtes : {nb_requetes}")
print(f"Succès : {succes}")
print(f"Durée : {duree:.1f}s")
print(f"Débit : {succes / duree:.1f} req/s")
Points clés à retenir
- Utilisez
AsyncOpenAIetasyncio.gatherpour la concurrence I/O - Limitez la concurrence avec un sémaphore pour respecter les quotas
- Le pattern map-reduce parallélise le traitement de longs documents
- Réservez
ProcessPoolExecutorau post-traitement CPU-intensif