Patterns avancés
Au-delà du workflow linéaire
Jusqu’ici, nous avons vu des workflows qui exécutent des activités de manière séquentielle. En pratique, les cas d’usage réels nécessitent des patterns plus sophistiqués : enchaîner des workflows, exécuter des activités en parallèle, ou prendre des décisions conditionnelles complexes.
Chaîner des workflows
Un workflow peut déclencher un autre workflow. C’est utile pour décomposer une logique complexe en sous-workflows réutilisables, ou pour contourner la limite de 51 200 événements par exécution.
import mistralai.workflows as workflows
from mistralai import Mistral
@workflows.activity()
async def lancer_sous_workflow(workflow_name: str, input_data: dict) -> dict:
"""Lance un workflow enfant et attend son résultat."""
client = Mistral()
execution = await client.workflows.execute_workflow_async(
workflow_identifier=workflow_name,
input=input_data
)
return execution.output
@workflows.workflow.define(name="pipeline_principal")
class PipelinePrincipal:
@workflows.workflow.entrypoint
async def run(self, params: dict) -> dict:
# Étape 1 : Collecte des données (sous-workflow dédié)
donnees = await lancer_sous_workflow(
"collecte_donnees",
{"sources": params["sources"]}
)
# Étape 2 : Analyse (sous-workflow dédié)
analyse = await lancer_sous_workflow(
"analyse_contenu",
{"donnees": donnees}
)
# Étape 3 : Génération du rapport (sous-workflow dédié)
rapport = await lancer_sous_workflow(
"generation_rapport",
{"analyse": analyse, "format": params["format"]}
)
return rapport
Pattern de continuation
Pour traiter un très grand nombre d’éléments, le workflow traite un lot puis se relance :
@workflows.workflow.define(name="traitement_par_lots")
class TraitementParLots:
@workflows.workflow.entrypoint
async def run(self, params: dict) -> dict:
offset = params.get("offset", 0)
taille_lot = 100
# Récupérer le lot courant
lot = await recuperer_lot(offset, taille_lot)
if not lot["items"]:
# Tous les lots sont traités
return {"status": "terminé", "total_traité": offset}
# Traiter le lot courant
await traiter_lot(lot["items"])
# Lancer le lot suivant via un nouveau workflow
suivant = await lancer_sous_workflow(
"traitement_par_lots",
{"offset": offset + taille_lot}
)
return suivant
Parallélisme avec asyncio.gather
Pour exécuter plusieurs activités en parallèle, utilisez asyncio.gather :
import asyncio
import mistralai.workflows as workflows
@workflows.activity()
async def analyser_source(url: str) -> dict:
"""Analyse une source de données."""
async with httpx.AsyncClient() as client:
response = await client.get(url)
return {"url": url, "taille": len(response.text)}
@workflows.workflow.define(name="analyse_parallele")
class AnalyseParallele:
@workflows.workflow.entrypoint
async def run(self, urls: list) -> dict:
# Lancer toutes les analyses en parallèle
resultats = await asyncio.gather(
*[analyser_source(url) for url in urls]
)
return {
"nombre_sources": len(resultats),
"resultats": list(resultats)
}
Avec asyncio.gather, les activités sont soumises simultanément à la plateforme. Cela réduit le temps total d’exécution par rapport à un traitement séquentiel.
Parallélisme avec gestion d’erreurs
Si certaines activités peuvent échouer sans bloquer les autres :
@workflows.workflow.define(name="analyse_tolerante")
class AnalyseToleranteAuxErreurs:
@workflows.workflow.entrypoint
async def run(self, urls: list) -> dict:
# return_exceptions=True empêche une erreur de bloquer les autres
resultats = await asyncio.gather(
*[analyser_source(url) for url in urls],
return_exceptions=True
)
succes = []
erreurs = []
for i, resultat in enumerate(resultats):
if isinstance(resultat, Exception):
erreurs.append({"url": urls[i], "erreur": str(resultat)})
else:
succes.append(resultat)
return {
"succes": succes,
"erreurs": erreurs,
"taux_succes": len(succes) / len(urls)
}
Workflows conditionnels
Les workflows peuvent contenir des branchements complexes basés sur les résultats des activités :
@workflows.workflow.define(name="moderation_contenu")
class ModerationContenu:
@workflows.workflow.entrypoint
async def run(self, contenu: dict) -> dict:
# Étape 1 : Analyse automatique
analyse = await analyser_contenu_auto(contenu["texte"])
if analyse["score_confiance"] > 0.95:
# Haute confiance : décision automatique
if analyse["categorie"] == "acceptable":
await publier_contenu(contenu["id"])
return {"decision": "publié", "methode": "automatique"}
else:
await rejeter_contenu(contenu["id"], analyse["raison"])
return {"decision": "rejeté", "methode": "automatique"}
elif analyse["score_confiance"] > 0.7:
# Confiance moyenne : deuxième analyse avec un autre modèle
seconde_analyse = await analyser_contenu_modele_b(contenu["texte"])
if analyse["categorie"] == seconde_analyse["categorie"]:
# Les deux modèles sont d'accord
if analyse["categorie"] == "acceptable":
await publier_contenu(contenu["id"])
return {"decision": "publié", "methode": "consensus_ia"}
else:
await rejeter_contenu(contenu["id"], analyse["raison"])
return {"decision": "rejeté", "methode": "consensus_ia"}
else:
# Désaccord : escalade humaine
await escalader_moderation(contenu["id"], analyse, seconde_analyse)
return {"decision": "en_attente", "methode": "escalade_humaine"}
else:
# Faible confiance : escalade directe
await escalader_moderation(contenu["id"], analyse, None)
return {"decision": "en_attente", "methode": "escalade_humaine"}
Pattern fan-out / fan-in
Ce pattern consiste à distribuer du travail en parallèle (fan-out) puis à agréger les résultats (fan-in) :
@workflows.workflow.define(name="analyse_multi_langue")
class AnalyseMultiLangue:
@workflows.workflow.entrypoint
async def run(self, document: dict) -> dict:
# Fan-out : traduire en parallèle dans 3 langues
traductions = await asyncio.gather(
traduire(document["texte"], "en"),
traduire(document["texte"], "de"),
traduire(document["texte"], "es")
)
# Fan-out : analyser chaque traduction en parallèle
analyses = await asyncio.gather(
*[analyser_sentiment(trad["texte_traduit"]) for trad in traductions]
)
# Fan-in : agréger les résultats
resume = await agreger_analyses(list(analyses))
return resume
Bonnes pratiques
- Décomposez les workflows complexes en sous-workflows réutilisables
- Parallélisez quand les activités sont indépendantes pour réduire le temps total
- Utilisez
return_exceptions=Truedansasyncio.gathersi certaines erreurs sont tolérables - Le pattern de continuation permet de contourner la limite d’événements pour les très longs traitements
- Gardez les branchements basés sur les résultats d’activités (pas de valeurs non déterministes)
Points clés à retenir
- Un workflow peut déclencher d’autres workflows via une activité qui appelle le SDK
- Le pattern de continuation permet de traiter un nombre illimité d’éléments par lots successifs
asyncio.gatherpermet d’exécuter des activités en parallèle pour réduire le temps total- Les branchements conditionnels complexes sont supportés tant qu’ils dépendent de résultats d’activités
- Le pattern fan-out / fan-in distribue le travail puis agrège les résultats