Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. Pode tentar alterar os diretórios.
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:
- Recolhe todas as mensagens pendentes do superstep anterior
- Encaminha mensagens para executores alvo com base em definições de arestas
- Executa todos os executores alvo de forma concorrente dentro do superpasso
- Espera que todos os executores completem antes de avançar (barreira de sincronização)
- 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:
- Executores — unidades de processamento num fluxo de trabalho
- Arestas — ligações entre executores
- Eventos — observabilidade do fluxo de trabalho
- Gestão de Estados