Créer une activité
Qu’est-ce qu’une activité ?
Une activité est l’unité opérationnelle de base dans un workflow. C’est là que le travail réel se passe : appels API externes, requêtes base de données, envoi d’emails, traitement de fichiers, appels aux modèles IA. Toute opération qui implique des entrées/sorties (I/O) ou des effets de bord doit être encapsulée dans une activité.
Contrairement au workflow qui orchestre, l’activité exécute. C’est la distinction fondamentale.
Anatomie d’une activité
Une activité est une fonction Python async décorée avec @workflows.activity() :
import mistralai.workflows as workflows
@workflows.activity()
async def recuperer_profil_utilisateur(user_id: str) -> dict:
"""Récupère le profil d'un utilisateur depuis l'API."""
async with httpx.AsyncClient() as client:
response = await client.get(f"https://api.example.com/users/{user_id}")
response.raise_for_status()
return response.json()
Analysons les éléments essentiels :
@workflows.activity()— Le décorateur qui enregistre la fonction comme activitéasync def— Obligatoire. Toutes les activités doivent être asynchrones- Paramètres JSON-sérialisables — Les arguments doivent être des types simples :
str,int,float,bool,list,dict, ou des modèles Pydantic - Retour JSON-sérialisable — Le résultat doit également être sérialisable en JSON
Exemples concrets d’activités
Appel à un modèle Mistral
import mistralai.workflows as workflows
from mistralai import Mistral
@workflows.activity()
async def generer_resume(texte: str) -> dict:
"""Génère un résumé du texte fourni via Mistral."""
client = Mistral()
response = await client.chat.complete_async(
model="mistral-large-latest",
messages=[
{"role": "system", "content": "Vous êtes un assistant spécialisé en résumé de texte."},
{"role": "user", "content": f"Résumez ce texte en 3 phrases :\n\n{texte}"}
]
)
return {
"resume": response.choices[0].message.content,
"tokens_utilises": response.usage.total_tokens
}
Requête base de données
import mistralai.workflows as workflows
import asyncpg
@workflows.activity()
async def sauvegarder_rapport(rapport: dict) -> dict:
"""Sauvegarde un rapport dans PostgreSQL."""
conn = await asyncpg.connect("postgresql://user:pass@localhost/db")
try:
row = await conn.fetchrow(
"INSERT INTO rapports (titre, contenu, date_creation) VALUES ($1, $2, NOW()) RETURNING id",
rapport["titre"],
rapport["contenu"]
)
return {"id": row["id"], "status": "sauvegardé"}
finally:
await conn.close()
Envoi d’email
import mistralai.workflows as workflows
import aiosmtplib
@workflows.activity()
async def envoyer_notification(destinataire: str, sujet: str, corps: str) -> dict:
"""Envoie un email de notification."""
message = f"Subject: {sujet}\nTo: {destinataire}\n\n{corps}"
await aiosmtplib.send(
message,
hostname="smtp.example.com",
port=587,
username="[email protected]",
password="secret",
use_tls=True
)
return {"status": "envoyé", "destinataire": destinataire}
Caractéristiques fondamentales
Pas de contrainte de déterminisme
Contrairement aux workflows, les activités n’ont aucune contrainte de déterminisme. Vous pouvez utiliser datetime.now(), uuid.uuid4(), random.random(), ou tout autre appel non déterministe librement dans une activité.
Unité atomique
Chaque appel à une activité est traité comme une unité atomique par la plateforme. L’activité soit réussit entièrement, soit échoue entièrement. Il n’y a pas d’état intermédiaire visible par le workflow.
Résultat enregistré
Une fois qu’une activité a réussi, son résultat est définitivement enregistré dans le journal d’événements du workflow. Lors d’un replay (reprise après panne), l’activité n’est jamais réexécutée — seul le résultat enregistré est utilisé.
Idempotence obligatoire
Comme une activité peut être retentée en cas d’échec (timeout réseau, erreur transitoire), elle doit être idempotente : exécutée plusieurs fois avec les mêmes paramètres, elle doit produire le même résultat.
@workflows.activity()
async def creer_ou_mettre_a_jour_client(email: str, nom: str) -> dict:
"""Idempotent : utilise un upsert au lieu d'un insert."""
async with httpx.AsyncClient() as client:
response = await client.put(
f"https://api.crm.com/clients/{email}",
json={"nom": nom, "email": email}
)
return response.json()
Nommer une activité
Vous pouvez donner un nom explicite à votre activité pour améliorer la lisibilité dans les traces :
@workflows.activity(name="Traitement des emails clients")
async def process_customer_emails(batch_id: str) -> dict:
# ...
pass
Ce nom apparaîtra dans les traces OpenTelemetry et facilitera le débogage.
Contraintes de taille
Les activités sont soumises à des limitations de taille :
- 2 MB maximum pour les paramètres d’entrée
- 2 MB maximum pour la valeur de retour
Si vous devez traiter des fichiers volumineux, stockez-les dans un service externe (S3, GCS) et passez uniquement l’URL ou l’identifiant dans l’activité.
Points clés à retenir
- Une activité est une fonction
asyncdécorée avec@workflows.activity()qui effectue le travail réel (I/O, API, BDD) - Les paramètres et le retour doivent être JSON-sérialisables
- Les activités n’ont pas de contrainte de déterminisme
- Une activité réussie n’est jamais réexécutée lors d’un replay
- Les activités doivent être idempotentes (même résultat si exécutées plusieurs fois)
- Limite de 2 MB en entrée et en sortie par activité