Generatore di flussi di lavoro ed esecuzione

Un flusso di lavoro collega executor e archi insieme in un grafico diretto e gestisce l'esecuzione. Coordina la chiamata dell'executor, il routing dei messaggi e lo streaming di eventi.

Creazione di flussi di lavoro

I flussi di lavoro vengono costruiti usando la WorkflowBuilder classe , che fornisce un'API Fluent per definire la struttura del flusso di lavoro:

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

I flussi di lavoro vengono costruiti usando 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()

workflow fornisce un modello di esecuzione basato su grafi in cui gli esecutori sono collegati da archi.

  • Executor : unità di elaborazione che riceve l'input e produce l'output
  • Edge : connette l'output di un executor all'input di un altro
  • Builder - Costruisce flussi di lavoro definendo esecutori e archi
  • Esecuzione : esegue un flusso di lavoro con l'input specificato
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
}

Esecuzione del flusso di lavoro

I flussi di lavoro supportano le modalità di esecuzione streaming e non-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()}")

Usare RunStreaming quando si desidera che si verifichino eventi:

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

Usare Run quando si vuole attendere il completamento del flusso di lavoro e quindi esaminare gli eventi raccolti:

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

È anche possibile esaminare gli eventi dell'executor raccolti da un'esecuzione non in streaming:

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

Suggerimento

Vedere gli esempi del flusso di lavoro per esempi eseguibili completi.

Convalida del flusso di lavoro

Il framework esegue la convalida completa durante la creazione di flussi di lavoro:

  • Compatibilità dei tipi: assicura che i tipi di messaggio siano compatibili tra executor connessi
  • Connettività del grafico: verifica che tutti gli executor siano raggiungibili dall'executor di avvio
  • Associazione dell'Executor: conferma che tutti gli executor sono associati e istanziati correttamente
  • Validazione dei bordi: Controlla i bordi duplicati e le connessioni non valide

Modello di esecuzione: passaggi sovrapposti

Il framework usa un modello di esecuzione Pregel modificato, ovvero un approccio BSP (Bulk Synchronous Parallel) con elaborazione basata su superstep.

Funzionamento di Supersteps

L'esecuzione del flusso di lavoro è organizzata in superstep distinti. Ogni passaggio superiore:

  1. Raccoglie tutti i messaggi in sospeso dal superstep precedente
  2. Indirizza i messaggi agli executor di destinazione in base alle definizioni di arco
  3. Esegue tutti gli esecutori di destinazione contemporaneamente all'interno del superstep
  4. Attende il completamento di tutti gli esecutori prima di avanzare (barriera di sincronizzazione)
  5. Accoda tutti i nuovi messaggi generati dagli executor per il passaggio successivo
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   │
└─────────────────┘    └─────────────────┘

Barriera di sincronizzazione

La caratteristica più importante è la barriera di sincronizzazione tra i superstep. All'interno di un unico passaggio, tutti gli executor attivati vengono eseguiti in parallelo, ma il flusso di lavoro non passa al passaggio successivo fino al completamento di ogni executor.

Ciò influisce sui modelli di distribuzione: se si distribuisce su più percorsi, uno con una catena di executor e un altro con un singolo executor a lunga esecuzione, il percorso concatenato non può avanzare fino al completamento dell'executor a lunga esecuzione.

Perché le superstep?

Il modello BSP offre garanzie importanti:

  • Esecuzione deterministica: dato lo stesso input, il flusso di lavoro viene sempre eseguito nello stesso ordine
  • Checkpoint affidabile: lo stato può essere salvato in limiti di passaggio superiore per la tolleranza di errore
  • Ragionamento più semplice: nessuna race condition tra superstep; ognuno vede una visualizzazione coerente dei messaggi

Uso del modello superstep

Se sono necessari percorsi paralleli veramente indipendenti che non si bloccano tra loro, consolidare i passaggi sequenziali in un singolo executor. Invece di concatenare step1 → step2 → step3, combina tale logica in un unico executor. Entrambi i percorsi paralleli vengono quindi eseguiti all'interno di un unico superstep.

Passaggi successivi

Argomenti correlati: