asyncio.gather pour traitement par lots
Orchestrer des requetes paralleles
asyncio.gather est la fonction centrale pour executer plusieurs coroutines en parallele et collecter tous les resultats. Combinee avec un semaphore, elle permet de traiter des lots de requetes de maniere efficace tout en respectant les limites de l’API.
Syntaxe de base
import asyncio
requests = ["Question 1", "Question 2", "Question 3"]
tasks = [process_request(r) for r in requests]
responses = await asyncio.gather(*tasks)
asyncio.gather prend un nombre variable de coroutines (d’ou le *tasks pour decompresser la liste) et retourne une liste de resultats dans le meme ordre que les taches soumises. C’est un point important : meme si la question 3 termine avant la question 1, responses[0] contiendra toujours le resultat de la question 1.
Pattern complet : lot avec semaphore
Voici le pattern que vous utiliserez le plus souvent en production :
from asyncio import Semaphore
from openai import AsyncOpenAI
import httpx
client = AsyncOpenAI(
api_key=os.getenv("XAI_API_KEY"),
base_url="https://api.x.ai/v1",
timeout=httpx.Timeout(3600.0)
)
async def process_batch(prompts: list[str], max_concurrent: int = 5) -> list[str]:
sem = Semaphore(max_concurrent)
async def single_request(prompt: str) -> str:
async with sem:
response = await client.chat.completions.create(
model="grok-4",
messages=[{"role": "user", "content": prompt}]
)
return response.choices[0].message.content
tasks = [single_request(p) for p in prompts]
return await asyncio.gather(*tasks)
Ce pattern encapsule le semaphore dans la fonction de lot. Vous n’avez qu’a appeler process_batch(mes_prompts) et recuperer la liste des reponses.
Gerer les erreurs dans un lot
Par defaut, si une seule coroutine leve une exception, asyncio.gather annule tout et propage l’erreur. Le parametre return_exceptions=True change ce comportement :
async def safe_batch(prompts: list[str]) -> list[str | Exception]:
sem = Semaphore(5)
async def single_request(prompt: str) -> str:
async with sem:
response = await client.chat.completions.create(
model="grok-4",
messages=[{"role": "user", "content": prompt}]
)
return response.choices[0].message.content
tasks = [single_request(p) for p in prompts]
results = await asyncio.gather(*tasks, return_exceptions=True)
return results
# Utilisation
results = await safe_batch(mes_prompts)
for i, result in enumerate(results):
if isinstance(result, Exception):
print(f"Erreur sur le prompt {i}: {result}")
else:
print(f"Resultat {i}: {result[:100]}...")
Avec return_exceptions=True, les erreurs sont capturees comme des valeurs dans la liste de resultats. Vous pouvez ensuite identifier et retraiter les requetes echouees.
Decouper en sous-lots (chunking)
Pour les tres gros volumes, il est preferable de decouper en sous-lots plutot que de lancer des milliers de taches simultanement :
async def chunked_batch(
prompts: list[str],
chunk_size: int = 50,
max_concurrent: int = 5
) -> list[str]:
all_results = []
for i in range(0, len(prompts), chunk_size):
chunk = prompts[i:i + chunk_size]
print(f"Traitement du lot {i//chunk_size + 1} ({len(chunk)} prompts)")
results = await process_batch(chunk, max_concurrent)
all_results.extend(results)
return all_results
Le chunking offre plusieurs avantages : vous pouvez afficher la progression, sauvegarder les resultats intermediaires, et limiter la memoire utilisee. Si un sous-lot echoue, vous ne perdez pas les resultats des sous-lots precedents.
Suivi de progression
Pour un lot de grande taille, un indicateur de progression est utile :
import asyncio
from asyncio import Semaphore
async def batch_with_progress(prompts: list[str], max_concurrent: int = 5) -> list[str]:
sem = Semaphore(max_concurrent)
completed = 0
total = len(prompts)
async def single_request(prompt: str) -> str:
nonlocal completed
async with sem:
response = await client.chat.completions.create(
model="grok-4",
messages=[{"role": "user", "content": prompt}]
)
completed += 1
print(f"\r{completed}/{total} termine", end="", flush=True)
return response.choices[0].message.content
tasks = [single_request(p) for p in prompts]
results = await asyncio.gather(*tasks)
print() # nouvelle ligne apres la progression
return results
Points cles a retenir
asyncio.gather(*tasks)execute des coroutines en parallele et retourne les resultats dans l’ordre originalreturn_exceptions=Trueempeche qu’une seule erreur ne fasse echouer tout le lot- Le chunking permet de gerer les tres gros volumes sans saturer la memoire
- Combinez toujours
gatheravec un semaphore pour respecter les rate limits - Ajoutez un suivi de progression pour les traitements de longue duree