工作流生成器和执行

工作流将执行器和边缘关联到定向图中,并管理执行。 它协调执行程序调用、消息路由和事件流式处理。

生成工作流

工作流是使用 WorkflowBuilder 类构造的,该类提供用于定义工作流结构的 Fluent API:

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

工作流是使用WorkflowBuilder类构造的。

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 包提供基于图形的执行模型,执行程序通过边缘进行连接。

  • 执行程序 - 接收输入和生成输出的处理单元
  • Edge - 将一个执行程序的输出连接到另一个执行程序的输入
  • 生成器 - 通过定义执行器和边缘来构造工作流
  • 运行 - 使用给定输入执行工作流
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
}

工作流执行

工作流支持流式处理和非流式处理执行模式:

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

如果您希望在事件发生时接收事件,请使用 RunStreaming

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

如果要等待工作流完成,然后检查收集的事件,请使用 Run

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

您还可以查看非流式运行收集的执行器事件:

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

Tip

有关完整的可运行示例,请参阅 工作流示例

工作流验证

生成工作流时,框架会执行全面的验证:

  • 类型兼容性:确保消息类型在连接的执行程序之间兼容
  • 图形连接:验证可从启动执行程序访问所有执行程序
  • 执行程序绑定:确认所有执行程序都已正确绑定和实例化
  • 边缘验证:检查重复边缘和无效连接

执行模型:超级步骤

该框架使用经过修改的 Pregel 执行模型 , 这是一种批量同步并行 (BSP) 方法,采用基于超步骤的处理。

Supersteps 的工作原理

工作流执行组织为离散的超级步骤。 每个超级步骤:

  1. 从上一个超级步骤收集所有挂起的消息
  2. 基于边缘定义将消息路由到目标执行程序
  3. 在超级步骤中并发运行所有目标执行程序
  4. 等待所有执行程序在推进之前完成(同步屏障)
  5. 将执行者发出的任何新消息排队处理,以用于下一个超级步骤
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   │
└─────────────────┘    └─────────────────┘

同步屏障

最重要的特征是超级步骤之间的同步屏障。 在单个超级步骤中,所有触发的执行程序都并行运行,但在每个执行程序完成之前,工作流不会前进到下一个超级步骤。

这会影响扇出模式:如果扇出到多个路径(一个具有执行程序链,另一个具有单个长时间运行的执行程序),则链式路径在长时间运行的执行程序完成之前无法前进。

为什么选择超级步骤?

BSP 模型提供重要保证:

  • 确定性执行:给定相同的输入,工作流始终按相同顺序执行
  • 可靠检查点:状态可以在超阶边界处保存,以实现容错
  • 更简单的推理:超级步骤之间没有竞争条件;每个超级步骤看到消息的一致性

使用 Superstep 模型

如果需要真正独立且互不阻塞的并行路径,可以将顺序步骤合并到单个执行程序中。 不要将step1 → step2 → step3逻辑链接起来,而是将其合并到一个单一的执行器中。 然后,两个并行路径在单个超级步骤中执行。

后续步骤

相关主题