Aller au contenu principal

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=True dans asyncio.gather si 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.gather permet 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