Async et bonnes pratiques
Mis à jour le 29 juillet 2026
Async I/O : une obligation technique
Toutes les activités dans les Workflows Mistral doivent utiliser des opérations d’entrée/sortie asynchrones. Ce n’est pas une recommandation de style, c’est une contrainte technique : un appel synchrone bloquant comme requests.post() immobilise le worker et empêche l’exécution de toute autre activité.
La raison tient à l’architecture du worker, qui exécute les activités dans une boucle d’événements asyncio. Quand une coroutine attend une réponse réseau de façon asynchrone, elle rend la main à la boucle, qui en profite pour faire avancer les autres activités. Un appel synchrone, lui, ne rend rien.
import requests
# INCORRECT : bloque le worker entier
@workflows.activity()
async def activite_bloquante(url: str) -> dict:
# requests.post() est synchrone — il bloque la boucle d'événements
response = requests.post(url, json={"query": "test"})
return response.json()
Le symptôme est déroutant en production : vous lancez dix activités en parallèle avec asyncio.gather, et elles s’exécutent en réalité les unes après les autres. Si chacune attend deux secondes une API distante, vous mesurez vingt secondes là où vous en attendiez deux, sans qu’aucune erreur ne soit remontée.
Choisir des librairies nativement asynchrones
La correction consiste à remplacer chaque librairie bloquante par son équivalent async. Pour les appels HTTP, httpx est le choix par défaut.
import httpx
# CORRECT : non bloquant
@workflows.activity()
async def activite_async(url: str) -> dict:
async with httpx.AsyncClient() as client:
response = await client.post(url, json={"query": "test"})
return response.json()
aiohttp remplit le même office avec une API légèrement différente, à base de sessions.
import aiohttp
@workflows.activity()
async def activite_aiohttp(url: str) -> dict:
async with aiohttp.ClientSession() as session:
async with session.post(url, json={"query": "test"}) as response:
return await response.json()
Côté base de données, asyncpg remplace psycopg2 pour PostgreSQL. Notez le try/finally : la connexion doit être fermée même si la requête lève une exception, faute de quoi les retries successifs épuiseront le pool de connexions.
import asyncpg
# CORRECT : driver PostgreSQL asynchrone
@workflows.activity()
async def requete_bdd_async(user_id: str) -> dict:
conn = await asyncpg.connect("postgresql://user:pass@localhost/db")
try:
row = await conn.fetchrow(
"SELECT * FROM users WHERE id = $1", user_id
)
return dict(row) if row else {"erreur": "utilisateur non trouvé"}
finally:
await conn.close()
L’accès aux fichiers suit la même logique avec aiofiles, y compris pour de petits fichiers où l’on serait tenté de garder open().
import aiofiles
# CORRECT : lecture de fichier asynchrone
@workflows.activity()
async def lire_fichier_async(chemin: str) -> dict:
async with aiofiles.open(chemin, mode="r") as f:
contenu = await f.read()
return {"contenu": contenu, "taille": len(contenu)}
Le SDK Mistral lui-même n’échappe pas à la règle : chaque méthode possède une variante suffixée _async, et c’est celle-là qu’il faut appeler depuis une activité.
from mistralai import Mistral
@workflows.activity()
async def appeler_mistral(prompt: str) -> dict:
client = Mistral()
# Utilisez complete_async, PAS complete
response = await client.chat.complete_async(
model="mistral-large-latest",
messages=[{"role": "user", "content": prompt}]
)
return {"reponse": response.choices[0].message.content}
Gardez ce tableau de correspondance sous la main lors de vos revues de code : chaque ligne de la colonne de gauche repérée dans une activité est un bug de performance en puissance.
| Opération | Incorrect (bloquant) | Correct (async) |
|---|---|---|
| HTTP | requests.get() | httpx.AsyncClient().get() |
| HTTP | urllib.request.urlopen() | aiohttp.ClientSession().get() |
| PostgreSQL | psycopg2.connect() | asyncpg.connect() |
| MySQL | mysql.connector.connect() | aiomysql.connect() |
| Redis | redis.Redis() | redis.asyncio.Redis() |
| Fichiers | open() | aiofiles.open() |
| Sleep | time.sleep() | asyncio.sleep() |
| Mistral | client.chat.complete() | client.chat.complete_async() |
Quand aucune version async n’existe
Certaines librairies métier, un SDK propriétaire ou un client legacy par exemple, n’offrent tout simplement pas de variante asynchrone. asyncio.to_thread() permet alors d’exécuter le code bloquant dans un thread séparé, sans immobiliser la boucle d’événements.
import asyncio
from some_library import blocking_function
@workflows.activity()
async def activite_compatibilite(params: dict) -> dict:
"""Encapsule un appel bloquant dans un thread."""
resultat = await asyncio.to_thread(
blocking_function,
params["arg1"],
params["arg2"]
)
return {"résultat": resultat}
C’est un compromis, pas une solution équivalente : chaque appel consomme un thread du pool, dont la taille est finie. Sur une activité appelée massivement en parallèle, vous déplacez le goulot d’étranglement au lieu de le supprimer. Préférez toujours une librairie nativement async quand elle existe.
Concevoir des activités propres
Au-delà de l’asynchronisme, quatre habitudes distinguent une activité robuste d’une activité qui vous coûtera des heures de débogage.
La première est le périmètre. Une activité porte une responsabilité et une seule, parce que c’est aussi l’unité de retry et l’unité de persistance : une activité fourre-tout qui échoue à l’envoi d’e-mail rejouera l’extraction et l’analyse depuis le début.
# BON : chaque activité a une responsabilité claire
@workflows.activity()
async def extraire_texte(url: str) -> dict:
# ...
pass
@workflows.activity()
async def analyser_sentiment(texte: str) -> dict:
# ...
pass
# MAUVAIS : activité fourre-tout
@workflows.activity()
async def tout_faire(url: str) -> dict:
texte = ... # extraction
sentiment = ... # analyse
# envoi email
# sauvegarde BDD
pass
La deuxième concerne les ressources. Le gestionnaire de contexte async with garantit la fermeture du client même en cas d’exception, ce qu’un appel manuel oublie tôt ou tard.
@workflows.activity()
async def requete_propre(params: dict) -> dict:
# Utilisez async with pour garantir la fermeture
async with httpx.AsyncClient() as client:
response = await client.get(params["url"])
# Le client est automatiquement fermé ici
return response.json()
La troisième est la plus subtile : toutes les erreurs ne méritent pas d’être rejouées. Un timeout, un 429 ou un 500 sont transitoires et doivent remonter pour déclencher un retry. Une erreur 4xx, en revanche, traduit une requête invalide que rejouer trois fois ne corrigera pas — autant retourner un résultat explicite et laisser le workflow décider.
@workflows.activity(
start_to_close_timeout=timedelta(minutes=2),
retry_policy_max_attempts=3,
)
async def activite_robuste(params: dict) -> dict:
async with httpx.AsyncClient(timeout=60) as client:
try:
response = await client.post(params["url"], json=params["data"])
response.raise_for_status()
return {"status": "succes", "data": response.json()}
except httpx.TimeoutException:
raise # Sera retried par la plateforme
except httpx.HTTPStatusError as e:
if e.response.status_code == 429:
raise # Rate limited — sera retried
elif e.response.status_code >= 500:
raise # Erreur serveur — sera retried
else:
# Erreur client (4xx) — ne pas retrier
return {"status": "erreur", "code": e.response.status_code}
La quatrième, enfin, découle de la limite de 2 MB par invocation : ne renvoyez au workflow que ce dont il a réellement besoin pour décider de la suite. Un décompte, un échantillon et l’URL du résultat complet suffisent presque toujours.
@workflows.activity()
async def extraire_donnees(source: str) -> dict:
donnees = await recuperer_grandes_donnees(source)
# Ne retournez que ce dont le workflow a besoin
return {
"nombre_elements": len(donnees),
"resume": donnees[:10], # Seulement les 10 premiers
"url_complete": await uploader_resultats(donnees)
}
Points clés à retenir
- Async I/O obligatoire : utilisez
httpx,asyncpg,aiofilesau lieu de leurs équivalents synchrones - Les appels bloquants (
requests,psycopg2,open) bloquent tout le worker - Utilisez
asyncio.to_thread()comme solution de dernier recours pour les librairies sans version async - Chaque activité doit avoir une seule responsabilité claire
- Fermez toujours les connexions avec
async with - Gérez les erreurs explicitement : certaines doivent être retried, d’autres non
- Limitez la taille des retours pour rester sous la limite de 2 MB