Criador e Execução de Fluxos de Trabalho

Um fluxo de trabalho liga executores e arestas num grafo direcionado e gere a execução. Coordena a invocação do executor, o encaminhamento 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 do 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();

Os 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 grafos onde os executores estão ligados por arestas.

  • Executor - Uma unidade de processamento que recebe entrada e produz saída
  • Edge - Liga a saída de um executor à entrada de outro
  • Builder - Constrói fluxos de trabalho definindo executores e arestas
  • Run - Executa um fluxo de trabalho com a entrada dada
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 do fluxo de trabalho

Os fluxos de trabalho suportam tanto os modos de execução em fluxo contínuo como os modos não em fluxo contínuo.

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

Utilize RunStreaming quando quiser os eventos à medida que acontecem:

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 esperar a conclusão do fluxo de trabalho e depois inspecionar os eventos recolhidos:

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

Também pode inspecionar os eventos do executor recolhidos durante 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)
    }
}

Tip

Consulte os exemplos de workflow para amostras completas executáveis.

Validação do fluxo de trabalho

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

  • Compatibilidade de tipo: garante que os tipos de mensagem sejam compatíveis entre executores conectados
  • Conectividade do gráfico: verifica se todos os executores estão acessíveis desde o executor inicial
  • Vinculação do Executor: Confirma que todos os executores estão corretamente vinculados e instanciados
  • Validação de Borda: Verifica se há bordas duplicadas e conexões inválidas

Modelo de Execução: Supersteps

O framework utiliza um modelo de execução Pregel modificado — uma abordagem Bulk Synchronous Parallel (BSP) com processamento baseado em superpassos.

Como Funcionam os Supersteps

A execução do fluxo de trabalho está organizada em superpassos discretos. Cada superpasso:

  1. Recolhe todas as mensagens pendentes do superstep anterior
  2. Encaminha mensagens para executores alvo com base em definições de arestas
  3. Executa todos os executores alvo de forma concorrente dentro do superpasso
  4. Espera que todos os executores completem antes de avançar (barreira de sincronização)
  5. Coloca em fila quaisquer novas mensagens emitidas pelos executores para o próximo 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   │
└─────────────────┘    └─────────────────┘

Barreira de Sincronização

A característica mais importante é a barreira de sincronização entre superetapas. Dentro de um único superpasso, todos os executores desencadeados correm em paralelo, mas o fluxo de trabalho não avança para o superpasso seguinte até que todos os executores terminem.

Isto afeta os padrões de dispersão: se se espalhar por múltiplos caminhos — um com uma cadeia de executores e outro com um único executor de longa duração — o caminho encadeado não pode avançar até que o executor de longa duração esteja concluído.

Porquê Supersteps?

O modelo BSP fornece garantias importantes:

  • Execução determinística: Dado o mesmo input, o fluxo de trabalho executa-se sempre na mesma ordem
  • Checkpointing fiável: O estado pode ser salvo nos limites dos superpassos para assegurar a tolerância a falhas.
  • Raciocínio mais simples: Sem condições de corrida entre superpassos; cada um tem uma visão consistente das mensagens

Trabalhar com o Modelo Superstep

Se precisares de caminhos paralelos verdadeiramente independentes que não se bloqueiem mutuamente, consolida passos sequenciais num único executor. Em vez de encadear step1 → step2 → step3, combine essa lógica num único executor. Ambos os caminhos paralelos executam-se então dentro de um único superpasso.

Passos seguintes

Tópicos relacionados: