Générateur de flux de travail & Exécution

Un flux de travail lie des exécuteurs et des arêtes pour former un graphe dirigé et gère l’exécution. Il coordonne l’appel de l’exécuteur, le routage des messages et le streaming d’événements.

Création de flux de travail

Les flux de travail sont construits à l’aide de la WorkflowBuilder classe, qui fournit une API Fluent pour définir la structure de flux de travail :

using Microsoft.Agents.AI.Workflows;

var processor = new DataProcessor();
var validator = new Validator();
var formatter = new Formatter();

// Build workflow
WorkflowBuilder builder = new(processor); // Set starting executor
builder.AddEdge(processor, validator);
builder.AddEdge(validator, formatter);
var workflow = builder.Build();

Les flux de travail sont construits à l’aide de la WorkflowBuilder classe :

from agent_framework import WorkflowBuilder

processor = DataProcessor()
validator = Validator()
formatter = Formatter()

# Build workflow
builder = WorkflowBuilder(start_executor=processor)
builder.add_edge(processor, validator)
builder.add_edge(validator, formatter)
workflow = builder.build()

Le workflow package fournit un modèle d’exécution basé sur des graphiques où les exécuteurs sont connectés par des arêtes.

  • Exécuteur : unité de traitement qui reçoit l’entrée et produit la sortie
  • Edge : connecte la sortie d’un exécuteur à l’entrée d’un autre
  • Générateur : construit des flux de travail en définissant des exécuteurs et des arêtes
  • Exécuter : exécute un flux de travail avec une entrée donnée
import (
    "github.com/microsoft/agent-framework-go/workflow"
    "github.com/microsoft/agent-framework-go/workflow/inproc"
)

uppercase := workflow.NewExecutor("UppercaseExecutor", func(input string) string {
    return strings.ToUpper(input)
}).Bind()

reverse := workflow.NewExecutor("ReverseExecutor", func(input string) string {
    runes := []rune(input)
    slices.Reverse(runes)
    return string(runes)
}).Bind()

wf, err := workflow.NewBuilder(uppercase).
    AddEdge(uppercase, reverse).
    WithOutputFrom(reverse).
    Build()
if err != nil {
    return err
}

Exécution du flux de travail

Les flux de travail prennent en charge les modes d’exécution de streaming et de non-diffusion en continu :

using Microsoft.Agents.AI.Workflows;

// Streaming execution — get events as they happen
StreamingRun run = await InProcessExecution.RunStreamingAsync(workflow, inputMessage);
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
    if (evt is ExecutorCompletedEvent executorComplete)
    {
        Console.WriteLine($"{executorComplete.ExecutorId}: {executorComplete.Data}");
    }

    if (evt is WorkflowOutputEvent outputEvt)
    {
        Console.WriteLine($"Workflow completed: {outputEvt.Data}");
    }
}

// Non-streaming execution — wait for completion
Run result = await InProcessExecution.RunAsync(workflow, inputMessage);
foreach (WorkflowEvent evt in result.NewEvents)
{
    if (evt is WorkflowOutputEvent outputEvt)
    {
        Console.WriteLine($"Final result: {outputEvt.Data}");
    }
}
# Streaming execution — get events as they happen
async for event in workflow.run(input_message, stream=True):
    if event.type == "output":
        print(f"Workflow completed: {event.data}")

# Non-streaming execution — wait for completion
events = await workflow.run(input_message)
print(f"Final result: {events.get_outputs()}")

Utilisez RunStreaming quand vous souhaitez que les événements se produisent :

stream, err := inproc.Default.RunStreaming(context.Background(), wf, "Hello, World!")
if err != nil {
    return err
}
defer stream.Close(context.Background())

for evt, err := range stream.WatchStream(context.Background()) {
    if err != nil {
        return err
    }
    if output, ok := evt.(workflow.OutputEvent); ok {
        fmt.Printf("Workflow completed: %v\n", output.Output)
    }
}

Utilisez Run quand vous souhaitez attendre la fin du flux de travail, puis inspectez les événements collectés :

run, err := inproc.Default.Run(context.Background(), wf, "Hello, World!")
if err != nil {
    return err
}

for evt := range run.NewEvents() {
    if output, ok := evt.(workflow.OutputEvent); ok {
        fmt.Printf("Final result: %v\n", output.Output)
    }
}

Vous pouvez également examiner les événements de l’exécuteur recueillis lors d’une exécution non streaming :

for evt := range run.NewEvents() {
    if evt, ok := evt.(workflow.ExecutorCompletedEvent); ok {
        fmt.Printf("%s: %v\n", evt.ExecutorID, evt.Result)
    }
}

Tip

Consultez les exemples de flux de travail pour obtenir des exemples exécutables complets.

Validation du flux de travail

L’infrastructure effectue une validation complète lors de la création de flux de travail :

  • Compatibilité des types de messages : garantit que les types de messages sont compatibles entre les exécuteurs connectés
  • Connectivité du graphique : vérifie que tous les exécuteurs sont accessibles à partir de l’exécuteur de démarrage
  • Liaison d’exécuteur : confirme que tous les exécuteurs sont correctement liés et instanciés
  • Validation des arêtes : vérifie les arêtes dupliquées et les connexions non valides

Modèle d’exécution : Supersteps

L’infrastructure utilise un modèle d’exécution Pregel modifié , une approche BSP (Bulk Synchronous Parallel) avec un traitement basé sur un superstep.

Fonctionnement des supersteps

L’exécution du flux de travail est organisée en supersteps discrets. Chaque superstep :

  1. Collecte tous les messages en attente à partir du superstep précédent
  2. Route les messages vers les exécuteurs cibles en fonction des définitions de périphérie
  3. Exécute tous les exécuteurs cibles simultanément dans le superstep
  4. Attendre que tous les exécuteurs soient terminés avant de progresser (barrière de synchronisation)
  5. Met en file d’attente tous les nouveaux messages émis par les exécuteurs pour le superstep suivant.
Superstep N:
┌─────────────────┐    ┌─────────────────┐    ┌─────────────────┐
│  Collect All    │───▶│  Route Messages │───▶│  Execute All    │
│  Pending        │    │  Based on Type  │    │  Target         │
│  Messages       │    │  & Conditions   │    │  Executors      │
└─────────────────┘    └─────────────────┘    └─────────────────┘
                                                       │
                                                       │ (barrier: wait for all)
┌─────────────────┐    ┌─────────────────┐             │
│  Start Next     │◀───│  Emit Events &  │◀────────────┘
│  Superstep      │    │  New Messages   │
└─────────────────┘    └─────────────────┘

Barrière de synchronisation

La caractéristique la plus importante est la barrière de synchronisation entre les supersteps. Dans un superstep unique, tous les exécuteurs déclenchés s’exécutent en parallèle, mais le workflow ne passe pas au superstep suivant tant que chaque exécuteur n’est pas terminé.

Cela affecte les modèles de fan-out : si vous divisez en plusieurs chemins (un avec une chaîne d’exécuteurs et un autre avec un seul exécuteur de longue durée), le chemin chaîné ne peut pas avancer tant que l’exécuteur de longue durée n’a pas terminé.

Pourquoi supersteps ?

Le modèle BSP fournit des garanties importantes :

  • Exécution déterministe : étant donné la même entrée, le flux de travail s’exécute toujours dans le même ordre
  • Point de contrôle fiable : l’état peut être enregistré aux limites de superétape pour la tolérance aux pannes
  • Raisonnement plus simple : Aucune condition de course entre les supersteps ; chacun voit une vue cohérente des messages

Utilisation du modèle Superstep

Si vous avez besoin de chemins parallèles réellement indépendants qui ne se bloquent pas, consolidez les étapes séquentielles dans un seul exécuteur. Au lieu de chaîner step1 → step2 → step3, combinez cette logique en un seul exécuteur. Les deux chemins parallèles s’exécutent ensuite au sein d’un seul superstep.

Étapes suivantes

Rubriques connexes :