Aller au contenu principal

Architecture asynchrone et files d'attente

Mis à jour le 28 juillet 2026

Pourquoi on ne bloque pas sur un appel LLM

Un appel API à un modèle de langage prend entre une et trente secondes selon le modèle et la complexité de la demande. Trente secondes, c’est une éternité pour un processus serveur : pendant ce temps, le thread qui attend ne fait rien, mais il occupe une place. Multipliez par cent requêtes simultanées et votre application tombe non pas parce qu’elle calcule trop, mais parce qu’elle attend trop.

Une architecture asynchrone adossée à une file d’attente répond à trois problèmes d’un coup. Elle absorbe les pics de charge, puisque la file grandit au lieu que les requêtes soient refusées. Elle centralise la logique de reprise, parce qu’une tâche qui échoue retourne dans la file au lieu de disparaître avec la requête HTTP qui l’a déclenchée. Et elle découple vos services : celui qui reçoit les demandes n’a plus besoin de connaître la durée de traitement de celui qui les exécute.

Le pattern producteur / consommateur

La structure de base tient en deux objets. Une tâche décrit ce qu’il faut faire et où en est le traitement ; une file distribue ces tâches à un pool de workers qui les consomment en parallèle. Modéliser le statut sous forme d’énumération plutôt que de chaîne libre évite la classe de bugs la plus banale de ce type de code : deux endroits du programme qui écrivent "terminee" et "terminée".

import asyncio
from dataclasses import dataclass, field
from enum import Enum

class StatutTache(Enum):
    EN_ATTENTE = "en_attente"
    EN_COURS = "en_cours"
    TERMINEE = "terminee"
    ECHOUEE = "echouee"

@dataclass
class Tache:
    id: str
    prompt: str
    modele: str = "gpt-5.6-terra"
    statut: StatutTache = StatutTache.EN_ATTENTE
    resultat: str | None = None
    tentatives: int = 0
    max_tentatives: int = 3

class FileTraitement:
    """File d'attente asynchrone pour les appels LLM."""

    def __init__(self, nb_workers: int = 5):
        self.file: asyncio.Queue[Tache] = asyncio.Queue()
        self.resultats: dict[str, Tache] = {}
        self.nb_workers = nb_workers

    async def ajouter(self, tache: Tache) -> str:
        """Ajoute une tâche à la file."""
        self.resultats[tache.id] = tache
        await self.file.put(tache)
        return tache.id

    async def worker(self, worker_id: int):
        """Worker qui traite les tâches de la file."""
        import openai
        client = openai.AsyncOpenAI()

        while True:
            tache = await self.file.get()
            tache.statut = StatutTache.EN_COURS
            tache.tentatives += 1

            try:
                response = await client.responses.create(
                    model=tache.modele,
                    input=tache.prompt,
                )
                tache.resultat = response.output_text
                tache.statut = StatutTache.TERMINEE

            except Exception as e:
                if tache.tentatives < tache.max_tentatives:
                    tache.statut = StatutTache.EN_ATTENTE
                    await self.file.put(tache)
                else:
                    tache.statut = StatutTache.ECHOUEE
                    tache.resultat = str(e)

            finally:
                self.file.task_done()

    async def demarrer(self):
        """Lance les workers."""
        workers = [
            asyncio.create_task(self.worker(i))
            for i in range(self.nb_workers)
        ]
        return workers

    async def attendre_fin(self):
        """Attend que toutes les tâches soient traitées."""
        await self.file.join()

Deux détails de ce worker méritent votre attention parce qu’ils reviennent dans toutes les implémentations réelles. Le compteur tentatives porté par la tâche elle-même, et non par le worker, permet à une tâche remise en file d’être reprise par n’importe quel autre worker sans perdre son historique. Le task_done() placé dans un bloc finally garantit que la file ne restera pas bloquée indéfiniment si une exception inattendue traverse le bloc except : sans lui, un join() qui n’aboutit jamais est le symptôme, et il est pénible à diagnostiquer.

Côté appelant, l’usage se lit sans surprise : on démarre les workers, on empile les tâches, on attend, on collecte, on arrête.

async def traiter_documents(documents: list[str]):
    file = FileTraitement(nb_workers=10)
    workers = await file.demarrer()

    # Ajouter les tâches
    for i, doc in enumerate(documents):
        tache = Tache(
            id=f"doc-{i}",
            prompt=f"Résumez ce document : {doc}",
        )
        await file.ajouter(tache)

    # Attendre la fin
    await file.attendre_fin()

    # Collecter les résultats
    for tache_id, tache in file.resultats.items():
        print(f"{tache_id}: {tache.statut.value}")

    # Arrêter les workers
    for w in workers:
        w.cancel()

Respecter les quotas plutôt que les subir

OpenAI impose deux limites simultanées : un nombre de requêtes par minute et un nombre de tokens par minute. Un pool de dix workers qui envoie tout ce qu’il peut atteindra l’une ou l’autre, recevra des erreurs de limitation, retentera, et aggravera la congestion. Il vaut mieux ralentir volontairement en amont que se faire refuser en aval.

Le limiteur ci-dessous applique un token bucket sur une fenêtre glissante d’une minute. Il tient deux historiques — les horodatages des requêtes récentes et les couples horodatage/tokens — les purge de ce qui est sorti de la fenêtre, puis attend si l’une des deux limites serait franchie. Le verrou asyncio est indispensable : sans lui, dix workers évalueraient la même situation en même temps et décideraient tous les dix qu’il reste de la place.

import asyncio
import time

class RateLimiter:
    """Limiteur de débit basé sur un token bucket."""

    def __init__(self, rpm: int = 500, tpm: int = 200_000):
        self.rpm = rpm
        self.tpm = tpm
        self.requetes_recentes: list[float] = []
        self.tokens_recents: list[tuple[float, int]] = []
        self.verrou = asyncio.Lock()

    async def attendre(self, tokens_estimes: int = 1000):
        """Attend si nécessaire pour respecter les limites."""
        async with self.verrou:
            maintenant = time.time()
            fenetre = 60.0

            # Nettoyer les entrées hors fenêtre
            self.requetes_recentes = [
                t for t in self.requetes_recentes
                if maintenant - t < fenetre
            ]
            self.tokens_recents = [
                (t, n) for t, n in self.tokens_recents
                if maintenant - t < fenetre
            ]

            # Vérifier RPM
            if len(self.requetes_recentes) >= self.rpm:
                attente = self.requetes_recentes[0] + fenetre - maintenant
                await asyncio.sleep(attente)

            # Vérifier TPM
            tokens_utilises = sum(n for _, n in self.tokens_recents)
            if tokens_utilises + tokens_estimes > self.tpm:
                attente = self.tokens_recents[0][0] + fenetre - maintenant
                await asyncio.sleep(attente)

            self.requetes_recentes.append(time.time())
            self.tokens_recents.append((time.time(), tokens_estimes))

Le paramètre tokens_estimes mérite un mot : vous ne connaissez pas le coût réel avant la réponse, donc vous travaillez sur une estimation. Prenez-la généreuse plutôt que juste. Sous-estimer systématiquement revient à désactiver la protection au moment précis où elle serait utile.

Rendre la main immédiatement

Pour les traitements de plusieurs minutes, faire patienter le client HTTP n’est plus tenable : les timeouts de proxy, de navigateur et de load balancer s’en mêleront avant vous. Le pattern webhook renverse la responsabilité. Votre endpoint accepte la demande, retourne aussitôt un identifiant de tâche, et l’appelant est notifié quand le résultat existe, sans avoir à interroger périodiquement.

from fastapi import FastAPI, BackgroundTasks

app = FastAPI()

@app.post("/api/analyse")
async def lancer_analyse(
    requete: dict,
    background_tasks: BackgroundTasks,
):
    """Lance une analyse en arrière-plan et notifie par webhook."""
    tache_id = generer_id()

    background_tasks.add_task(
        executer_analyse,
        tache_id=tache_id,
        document=requete["document"],
        webhook_url=requete.get("webhook_url"),
    )

    return {"tache_id": tache_id, "statut": "en_cours"}


async def executer_analyse(
    tache_id: str,
    document: str,
    webhook_url: str | None,
):
    """Exécute l'analyse et envoie le résultat par webhook."""
    import httpx

    resultat = await appeler_llm(document)

    if webhook_url:
        async with httpx.AsyncClient() as http:
            await http.post(webhook_url, json={
                "tache_id": tache_id,
                "statut": "terminee",
                "résultat": resultat,
            })

Conservez toujours l’identifiant retourné et l’état associé quelque part de durable. Un webhook peut se perdre, et l’appelant doit garder la possibilité d’aller chercher lui-même le résultat.

Points clés à retenir

  • Utilisez des files d’attente asynchrones pour découpler l’ingestion du traitement
  • Implémentez un rate limiter pour respecter les quotas OpenAI
  • Gérez les retries avec backoff exponentiel dans vos workers
  • Pour les tâches longues, préférez un pattern webhook plutôt que du polling