Generador de flujos de trabajo y ejecución

Un flujo de trabajo vincula los ejecutores y los bordes en un grafo dirigido y administra la ejecución. Coordina la invocación del ejecutor, el enrutamiento de mensajes y el streaming de eventos.

Creación de flujos de trabajo

Los flujos de trabajo se construyen mediante la WorkflowBuilder clase , que proporciona una API fluida para definir la estructura del flujo de trabajo:

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();

Los flujos de trabajo se construyen mediante la WorkflowBuilder clase :

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()

El workflow paquete proporciona un modelo de ejecución basado en grafos donde los ejecutores están conectados por bordes.

  • Ejecutor : una unidad de procesamiento que recibe la entrada y genera la salida.
  • Edge : conecta la salida de un ejecutor a la entrada de otra.
  • Generador : construye flujos de trabajo mediante la definición de ejecutores y bordes
  • Ejecutar : ejecuta un flujo de trabajo con una entrada determinada.
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
}

Ejecución del flujo de trabajo

Los flujos de trabajo admiten los modos de ejecución de streaming y no streaming:

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()}")

Use RunStreaming cuando desee eventos a medida que se produzcan:

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

Use Run cuando quiera esperar a que finalice el flujo de trabajo y, a continuación, inspeccione los eventos recopilados:

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

También puede inspeccionar los eventos del ejecutor recopilados en una ejecución sin streaming:

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

Tip

Consulte los ejemplos de flujo de trabajo para ver ejemplos ejecutables completos.

Validación de flujo de trabajo

El marco realiza una validación completa al compilar flujos de trabajo:

  • Compatibilidad de tipos: garantiza que los tipos de mensaje son compatibles entre los ejecutores conectados.
  • Conectividad de grafos: comprueba que todos los ejecutores son accesibles desde el ejecutor de inicio.
  • Enlace de ejecutores: confirma que todos los ejecutores están correctamente vinculados e instanciados
  • Validación de bordes: verifica bordes duplicados y conexiones no válidas

Modelo de ejecución: Supersteps

El marco utiliza un modelo de ejecución Pregel modificado: un enfoque de paralelismo sincrónico masivo (BSP) con procesamiento basado en superpasos.

Cómo funcionan los superpasos

La ejecución del flujo de trabajo se organiza en superpasos discretos. Cada superpaso:

  1. Recopila todos los mensajes pendientes del superpaso anterior.
  2. Enruta los mensajes a los ejecutores de destino en función de las definiciones perimetrales
  3. Ejecuta todos los ejecutores de destino simultáneamente dentro del superstep
  4. Espera a que todos los ejecutores completen su tarea antes de avanzar (barrera de sincronización)
  5. Pone en cola cualquier mensaje nuevo emitido por los ejecutores para el siguiente superpaso
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   │
└─────────────────┘    └─────────────────┘

Barrera de sincronización

La característica más importante es la barrera de sincronización entre supersteps. Dentro de un mismo superpaso, todos los ejecutores activados se ejecutan en paralelo, pero el flujo de trabajo no avanza al siguiente superpaso hasta que todos los ejecutores hayan finalizado.

Esto afecta a los patrones de ramificación: si se ramifica hacia múltiples rutas —una con una cadena de ejecutores y otra con un único ejecutor de larga duración—, la ruta encadenada no puede avanzar hasta que el ejecutor de larga duración haya finalizado.

¿Por qué supersteps?

El modelo BSP proporciona garantías importantes:

  • Ejecución determinista: dada la misma entrada, el flujo de trabajo siempre se ejecuta en el mismo orden.
  • Puntos de control fiables: el estado se puede guardar en los límites de los superpasos para garantizar la tolerancia a errores
  • Razonamiento más sencillo: no hay condiciones de carrera entre superpasos; cada uno ve una visión coherente de los mensajes

Trabajar con el modelo de Superstep

Si necesita rutas paralelas realmente independientes que no se bloqueen entre sí, consolide los pasos secuenciales en un único ejecutor. En lugar de encadenar step1 → step2 → step3, combina esa lógica en un único ejecutor. Después, ambas rutas paralelas se ejecutan dentro de un único superpaso.

Pasos siguientes

Temas relacionados: