Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
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 del ejecutor: confirma que todos los ejecutores están correctamente enlazados e instanciados.
- Validación de bordes: verifica bordes duplicados y conexiones no válidas
Modelo de ejecución: Supersteps
El framework utiliza un modelo de ejecución modificado de "Pregel" — un enfoque paralelo síncrono 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:
- Recopila todos los mensajes pendientes del superpaso anterior.
- Enruta los mensajes a los ejecutores de destino en función de las definiciones perimetrales
- Ejecuta todos los ejecutores de destino simultáneamente dentro del superstep
- Espera a que todas las tareas se hayan completado antes de avanzar (barrera de sincronización)
- Pone en cola los nuevos mensajes emitidos por ejecutores para el siguiente superstep
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 único superetapa, todos los ejecutores desencadenados se ejecutan en paralelo, pero el flujo de trabajo no avanza a la siguiente superetapa hasta que todos los ejecutores finalicen.
Esto afecta a los patrones de ramificación: si se divide en varias rutas (una con una cadena de ejecutores y otra con un único ejecutor de largo plazo), la ruta encadenada no puede avanzar hasta que se complete el ejecutor de largo plazo.
¿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.
- Punto de control confiable: el estado se puede guardar en límites de superpaso para la tolerancia a errores
- Razonamiento más sencillo: No hay condiciones de carrera entre superpasos; cada uno ve una vista 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, combine esa lógica en un ejecutor. Después, ambas rutas paralelas se ejecutan dentro de un único superpaso.
Pasos siguientes
Temas relacionados:
- Ejecutores : unidades de procesamiento en un flujo de trabajo
- Bordes : conexiones entre ejecutores
- Eventos : observabilidad del flujo de trabajo
- Administración de estados