Aller au contenu principal

Async et bonnes pratiques

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, c’est une contrainte technique : utiliser des appels synchrones bloquants (comme requests.post()) bloque le worker et empêche l’exécution d’autres activités.

Le problème du blocking I/O

Un worker exécute les activités dans une boucle d’événements asyncio. Lorsqu’une activité fait un appel réseau synchrone, elle bloque toute la boucle :

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()

Pendant que requests.post() attend la réponse du serveur (potentiellement plusieurs secondes), le worker ne peut exécuter aucune autre activité. Avec 10 activités en parallèle, elles s’exécutent toutes séquentiellement au lieu de simultanément.

La solution : async I/O systématique

Utilisez des librairies asynchrones pour toutes vos opérations I/O :

Appels HTTP : httpx (pas requests)

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()

Appels HTTP : aiohttp (alternative)

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()

Base de données : asyncpg (pas psycopg2)

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()

Fichiers : aiofiles (pas 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)}

SDK Mistral : utiliser la version async

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}

Tableau récapitulatif : correct vs incorrect

OpérationIncorrect (bloquant)Correct (async)
HTTPrequests.get()httpx.AsyncClient().get()
HTTPurllib.request.urlopen()aiohttp.ClientSession().get()
PostgreSQLpsycopg2.connect()asyncpg.connect()
MySQLmysql.connector.connect()aiomysql.connect()
Redisredis.Redis()redis.asyncio.Redis()
Fichiersopen()aiofiles.open()
Sleeptime.sleep()asyncio.sleep()
Mistralclient.chat.complete()client.chat.complete_async()

Gestion des librairies sans version async

Certaines librairies n’ont pas de version asynchrone. Dans ce cas, utilisez asyncio.to_thread() pour exécuter le code bloquant dans un thread séparé :

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 {"resultat": resultat}

Cette approche est un compromis : elle ne bloque pas la boucle d’événements, mais elle consomme un thread du pool. Préférez toujours une librairie nativement async quand elle existe.

Bonnes pratiques de conception des activités

1. Une responsabilité par activité

# 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

2. Toujours fermer les connexions

@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()

3. Gérer les erreurs explicitement

@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}

4. Limiter la taille des retours

@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, aiofiles au 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