Construção e Execução de Fluxo de Trabalho

Um fluxo de trabalho vincula executores e bordas em um grafo direcionado e gerencia a execução. Ele coordena a invocação do executor, o roteamento de mensagens e o streaming de eventos.

Criando fluxos de trabalho

Os fluxos de trabalho são construídos usando a WorkflowBuilder classe, que fornece uma API fluente para definir a estrutura de fluxo de trabalho:

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

Fluxos de trabalho são construídos usando a 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()

O workflow pacote fornece um modelo de execução baseado em grafo em que os executores são conectados por bordas.

  • Executor – Uma unidade de processamento que recebe entrada e produz saída
  • Borda – Conecta a saída de um executor à entrada de outro
  • Construtor – Constrói fluxos de trabalho definindo executores e bordas
  • Executar – Executa um fluxo de trabalho com determinada entrada
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
}

Execução de fluxo de trabalho

Os fluxos de trabalho dão suporte a modos de execução de streaming e não 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 quando quiser que os eventos ocorram:

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 quando quiser aguardar a conclusão do fluxo de trabalho e inspecione os eventos coletados:

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

Você também pode inspecionar os eventos do executor coletados em uma execução sem streaming:

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

Dica

Consulte os exemplos de fluxo de trabalho para obter exemplos executáveis completos.

Validação de fluxo de trabalho

A estrutura executa uma validação abrangente ao criar fluxos de trabalho:

  • Compatibilidade de tipos: garante que os tipos de mensagem sejam compatíveis entre executores conectados
  • Conectividade do Graph: verifica se todos os executores podem ser acessados no executor inicial
  • Vinculação do Executor: Confirma se todos os executores estão devidamente vinculados e instanciados
  • Validação de borda: verifica se há bordas duplicadas e conexões inválidas

Modelo de execução: Supersteps

A estrutura usa um modelo de execução de Pregel modificado – uma abordagem BSP (Bulk Synchronous Parallel) com processamento baseado em superstep.

Como funcionam os Supersteps

A execução do fluxo de trabalho é organizada em superpassos discretos. Cada superstep:

  1. Coleta todas as mensagens pendentes do superpasso anterior
  2. Roteia mensagens para executores de destino com base em definições de borda
  3. Executa todos os executores de destino simultaneamente dentro do superstep
  4. Aguarda que todos os executores sejam concluídos antes de avançar (barreira de sincronização)
  5. Enfileira todas as novas mensagens emitidas pelos executores para a próxima superetapa
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   │
└─────────────────┘    └─────────────────┘

Barreira de sincronização

A característica mais importante é a barreira de sincronização entre superpassos. Em um único superstep, todos os executores disparados são executados em paralelo, mas o fluxo de trabalho não avança para o próximo superstep até que todos os executores tenham sido concluídos.

Isso afeta os padrões de fan-out: se você usar vários caminhos , um com uma cadeia de executores e outro com um único executor de longa execução, o caminho encadeado não poderá avançar até que o executor de longa execução seja concluído.

Por que supersteps?

O modelo BSP fornece garantias importantes:

  • Execução determinística: dada a mesma entrada, o fluxo de trabalho sempre é executado na mesma ordem
  • Ponto de verificação confiável: o estado pode ser salvo em superetapas para tolerância a falhas
  • Raciocínio mais simples: sem condições de corrida entre superpassos; cada um vê uma visualização consistente de mensagens

Trabalhando com o Modelo Superstep

Se você precisar de caminhos paralelos verdadeiramente independentes que não se bloqueiem, consolide as etapas sequenciais em um único executor. Em vez de encadear step1 → step2 → step3, combine essa lógica em um executor. Os dois caminhos paralelos então são executados em uma única superetapa.

Próximas Etapas 

Tópicos relacionados: