Aller au contenu principal

Patterns avancés

Mis à jour le 29 juillet 2026

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. Ces quatre patterns — chaînage, continuation, parallélisme, fan-out/fan-in — couvrent la quasi-totalité des architectures que vous aurez à construire.

Chaîner des workflows

Un workflow peut en déclencher un autre. Cette capacité sert deux objectifs distincts : décomposer une logique complexe en sous-workflows réutilisables, et contourner la limite de 51 200 événements par exécution puisque chaque workflow enfant dispose de son propre budget. Le mécanisme passe par une activité, seul endroit où le SDK peut légitimement effectuer un appel réseau.

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",
            {"données": 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

L’intérêt dépasse la simple lisibilité : collecte_donnees devient un workflow autonome, testable seul et réutilisable par d’autres pipelines, avec ses propres timeouts et ses propres politiques de retry.

Le même mécanisme, appliqué à un workflow qui s’appelle lui-même, donne le pattern de continuation. Plutôt que de parcourir cent mille éléments dans une seule exécution — ce que le plafond d’événements interdit — le workflow traite un lot et délègue la suite à une nouvelle instance.

@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

Chaque nouvelle exécution repart avec un journal vierge. Le nombre d’éléments traitables devient donc illimité, la condition d’arrêt étant simplement un lot vide.

Paralléliser avec asyncio.gather

Lorsque plusieurs activités sont indépendantes les unes des autres, rien ne justifie de les exécuter en file. asyncio.gather les soumet simultanément à la plateforme, et le temps total se rapproche de celui de l’activité la plus lente au lieu d’être leur somme.

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),
            "résultats": list(resultats)
        }

Une réserve s’impose cependant : par défaut, gather propage la première exception rencontrée et abandonne l’ensemble. Si vous analysez quarante sources externes, il suffit qu’une seule soit hors ligne pour perdre les trente-neuf autres. Le paramètre return_exceptions=True change ce comportement : les erreurs sont renvoyées comme des valeurs ordinaires, à vous de les trier.

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

Ce workflow rend un résultat partiel accompagné d’un taux de succès, ce qui permet en aval de décider si la qualité est suffisante ou s’il faut relancer.

Décider dans le workflow

Les branchements complexes sont non seulement autorisés mais souvent l’intérêt principal du workflow, à condition que les conditions dérivent de résultats d’activités. Le cas de la modération de contenu illustre bien cette logique en cascade : décision automatique quand le modèle est très sûr de lui, second avis quand la confiance est moyenne, escalade humaine en cas de désaccord ou de faible confiance.

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

Remarquez le champ methode renvoyé à chaque sortie : il permet de mesurer a posteriori quelle proportion des contenus a été traitée sans intervention humaine, et donc d’ajuster les seuils.

Fan-out / fan-in

Le dernier pattern combine les précédents : distribuer le travail en parallèle, puis agréger les résultats en un point unique. Ici, un document est traduit simultanément en trois langues, chaque traduction est analysée en parallèle, et l’ensemble converge vers une synthèse.

@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

Les deux vagues de parallélisme sont successives, car la seconde dépend des résultats de la première. En séquentiel strict, ce workflow exigerait sept appels enchaînés ; ici, il en coûte trois paliers d’orchestration.

En résumé, décomposez les workflows complexes en sous-workflows réutilisables, parallélisez dès que les activités sont indépendantes, activez return_exceptions=True quand des échecs partiels sont tolérables, gardez la continuation en réserve pour les très longs traitements, et veillez toujours à ce que vos branchements reposent sur des résultats d’activités et non sur des 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