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