Aller au contenu principal

asyncio.gather pour traitement par lots

Mis à jour le 30 juillet 2026

Orchestrer des requêtes parallèles

asyncio.gather est la fonction centrale pour exécuter plusieurs coroutines en parallèle et collecter tous les résultats. Combinée avec un sémaphore, elle permet de traiter des lots de requêtes de manière 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’où le *tasks pour décompresser la liste) et retourne une liste de résultats dans le même ordre que les tâches soumises. C’est un point important : même si la question 3 se termine avant la question 1, responses[0] contiendra toujours le résultat de la question 1.

Pattern complet : lot avec sémaphore

Voici le pattern que vous utiliserez le plus souvent en production :

import os
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.5",
                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 sémaphore dans la fonction de lot. Vous n’avez qu’à appeler process_batch(mes_prompts) et récupérer la liste des réponses.

Gérer les erreurs dans un lot

Par défaut, si une seule coroutine lève une exception, asyncio.gather annule tout et propage l’erreur. Le paramètre 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.5",
                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"Résultat {i}: {result[:100]}...")

Avec return_exceptions=True, les erreurs sont capturées comme des valeurs dans la liste de résultats. Vous pouvez ensuite identifier et retraiter les requêtes échouées.

Découper en sous-lots (chunking)

Pour les très gros volumes, il est préférable de découper en sous-lots plutôt que de lancer des milliers de tâches simultanément :

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 résultats intermédiaires, et limiter la mémoire utilisée. Si un sous-lot échoue, vous ne perdez pas les résultats des sous-lots précédents.

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.5",
                messages=[{"role": "user", "content": prompt}]
            )
            completed += 1
            print(f"\r{completed}/{total} terminé", 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 après la progression
    return results

Points clés à retenir

  • asyncio.gather(*tasks) exécute des coroutines en parallèle et retourne les résultats dans l’ordre original
  • return_exceptions=True empêche qu’une seule erreur ne fasse échouer tout le lot
  • Le chunking permet de gérer les très gros volumes sans saturer la mémoire
  • Combinez toujours gather avec un sémaphore pour respecter les rate limits
  • Ajoutez un suivi de progression pour les traitements de longue durée