Workflow 进阶

上一篇介绍了 Executor、Edge 和 WorkflowBuilder,也解释了消息如何在工作流中流动。不过,当时的例子都是节点 A → 节点 B → 节点 C 这样的直线流程,没有分支,也不需要外部介入。

实际项目往往没有这么简单。例如,邮件需要按照 “垃圾、正常、不确定” 进入不同的处理流程;长时间运行的任务可能需要暂停,等待用户确认后再继续;多个 Agent 的结果需要全部返回后才能汇总;流程意外中断时,还要能从先前的状态恢复。只使用 AddEdge 把节点依次连起来,无法解决这些问题。

这一篇继续介绍 Workflow 的进阶能力,包括条件路由、并行与汇聚、有状态 Executor、事件、人工介入、检查点、子工作流,以及 SDK 内置的多 Agent 编排方式。


路由、分发与汇聚

上一篇只使用了直连边 AddEdge(source, target)。除此之外,SDK 还提供条件边、Switch、FanOut 和 FanIn,用来处理分支、并行和汇聚。

边的类型API干什么
直连AddEdge(s, t)无条件的线性管道
条件AddEdge<T>(s, t, condition)使用谓词决定是否进入目标节点
Switch-CaseAddSwitch(s, b => b.AddCase(...).WithDefault(...))按顺序匹配多路分支,可以设置默认分支
FanOutAddFanOutEdge(s, targets, targetSelector)将消息发送给一个或多个目标
FanInAddFanInBarrierEdge(sources, t)等待多个来源均有产出后再汇聚

除了直接传入 Executor,边的两端也可以只写字符串 ID。字符串会先转换成 ExecutorPlaceholder,因此可以先从配置中建立拓扑,稍后再通过 BindExecutor 绑定具体实现。

下面的流程由两条字符串定义组成。transform 节点使用哪个处理器,不写死在拓扑中,而是在应用启动时根据配置决定:

internal sealed class TextExecutor(string id, Func<string, string> transform)
    : Executor<string, string>(id)
{
    public override ValueTask<string> HandleAsync(
        string message,
        IWorkflowContext context,
        CancellationToken cancellationToken = default)
    {
        return ValueTask.FromResult(transform(message));
    }
}

string[] edgeDefinitions =
[
    "input -> transform",
    "transform -> output"
];

// 先解析字符串,只使用节点 ID 建立拓扑。
(string Source, string Target)[] edges = edgeDefinitions
    .Select(definition => definition.Split("->", StringSplitOptions.TrimEntries))
    .Select(parts => parts.Length == 2
        ? (Source: parts[0], Target: parts[1])
        : throw new InvalidOperationException($"无效的 Edge 定义: {string.Join(" -> ", parts)}"))
    .ToArray();

WorkflowBuilder builder = new(edges[0].Source);
foreach ((string source, string target) in edges)
{
    builder.AddEdge(source, target);
}
builder.WithOutputFrom("output");

// 运行前读取配置,选择 transform ID 对应的具体实现。
string processor = configuration["Workflow:Processor"] ?? "uppercase";
Func<string, string> transform = processor switch
{
    "uppercase" => text => text.ToUpperInvariant(),
    "trim" => text => text.Trim(),
    _ => throw new InvalidOperationException($"未知的处理器: {processor}")
};

builder.BindExecutor(new TextExecutor("input", text => text));
builder.BindExecutor(new TextExecutor("transform", transform));
builder.BindExecutor(new TextExecutor("output", text => text));

Workflow workflow = builder.Build();


这里,Edge 只依赖 inputtransformoutput 三个稳定 ID。不同环境可以为 transform 选择不同实现,也可以从数据库或配置文件读取整组边,再按 ID 注册相应的 Executor。这样改变执行器实现时不需要改动拓扑代码。

需要注意,所有 placeholder 都必须在 Build() 前完成绑定,否则构建时会抛出 InvalidOperationExceptionBuild() 得到的 Workflow 仍然不可变,不能在一次运行已经开始后替换节点;配置发生变化时,应重新创建 builder、完成绑定并构建新的 Workflow。


下面仍以邮件处理为例:垃圾检测 agent 分析邮件,然后按照检测结果进入不同的处理分支。

条件边适合处理简单分支。

一个源节点可以连接多条条件边,每条边的条件会被独立判断。

// 检测结果:is_spam + reason
public sealed class DetectionResult
{
    public bool IsSpam { get; set; }
    public string Reason { get; set; } = "";
}

// 条件工厂:复用同一种判断模式
private static Func<object?, bool> GetCondition(bool expected) =>
    msg => msg is DetectionResult r && r.IsSpam == expected;

var workflow = new WorkflowBuilder(spamDetectionExecutor)
    // 非垃圾 → 走邮件助手
    .AddEdge(spamDetectionExecutor, emailAssistantExecutor, condition: GetCondition(false))
    // 垃圾 → 走垃圾处理器
    .AddEdge(spamDetectionExecutor, handleSpamExecutor, condition: GetCondition(true))
    .WithOutputFrom(emailAssistantExecutor, handleSpamExecutor)
    .Build();


条件是一个 Func<T?, bool>,参数为源节点产出的消息,返回值表示是否沿这条边继续传递。需要注意,同一源节点上的多条条件边会独立判断,并不具有 if/else 的互斥关系。如果两个条件都返回 true,消息会同时进入两个目标节点。上面的例子通过 == true== false 保证两条路径互斥。


如果所有条件都不满足,消息不会继续传递。条件较多时,可以改用 Switch,并通过 WithDefault 设置默认路径。


Switch 将同一源节点的多路分支集中在一个配置块中,比散落的多条条件边更容易阅读,也支持默认分支:

public enum SpamDecision { NotSpam, Spam, Uncertain }

private static Func<object?, bool> GetCondition(SpamDecision expected) =>
    msg => msg is DetectionResult r && r.Decision == expected;

WorkflowBuilder builder = new(spamDetectionExecutor);
builder.AddSwitch(spamDetectionExecutor, sb => sb
    .AddCase(GetCondition(SpamDecision.NotSpam), emailAssistantExecutor)
    .AddCase(GetCondition(SpamDecision.Spam),    handleSpamExecutor)
    .WithDefault(handleUncertainExecutor)   // 兜底:前面都不匹配走这里
);
builder.AddEdge(emailAssistantExecutor, sendEmailExecutor);
builder.WithOutputFrom(handleSpamExecutor, sendEmailExecutor, handleUncertainExecutor);

Workflow workflow = builder.Build();


AddCase 会按照添加顺序逐个判断,命中第一个 case 后便停止;如果全部不匹配,则进入 WithDefault 指定的节点。因此,它对应的是 if/else if/else 语义。实现上,SwitchBuilder 最终会将这些配置转换成一条带目标选择器的 FanOut 边。

两路互斥分支通常使用条件边就够了;分支较多且只应选择其中一路时,Switch 会更清楚。


FanOut 用于将一条消息发送给多个目标。没有目标选择器时,消息会广播给所有目标:

// 一条邮件同时发给:回复助手、归档、分析
builder.AddFanOutEdge(emailAnalysisExecutor, [replyAssistant, archiver, analyzer]);


传入 targetSelector 后,还可以根据消息内容动态选择一个或多个目标:

// 目标顺序决定索引:[垃圾处理, 邮件助手, 摘要器, 不确定处理]
builder.AddFanOutEdge(
    emailAnalysisExecutor,
    targets: [handleSpam, emailAssistant, summarizer, handleUncertain],
    targetSelector: GetTargetAssigner());

// 选择器:返回"要激活的目标索引列表"
private static Func<AnalysisResult?, int, IEnumerable<int>> GetTargetAssigner() =>
    (result, targetCount) =>
    {
        if (result!.Decision == SpamDecision.Spam)     return [0];              // 只垃圾处理
        if (result.Decision == SpamDecision.Uncertain) return [3];              // 只不确定处理
        // NotSpam:总是给邮件助手;长邮件额外给摘要器
        List<int> targets = [1];
        if (result.EmailLength > 100) targets.Add(2);
        return targets;
    };


选择器的类型是 Func<T?, int, IEnumerable<int>>。第二个参数表示目标总数,返回值是需要激活的目标在 targets 列表中的索引。上面的代码中,普通长邮件会同时进入邮件助手和摘要器两个节点。


FanIn 用于汇聚多个来源。只有各个源节点都至少产生一条消息后,消息才会被送到目标节点:

// start 先 fan-out 到三个 worker,三者全部跑完,结果才汇聚到 aggregator
builder.AddFanOutEdge(start, [worker1, worker2, worker3]);
builder.AddFanInBarrierEdge(sources: [worker1, worker2, worker3], target: aggregator);


这里的屏障判断的是每个源都有产出,而不是各个源的消息数量相同。如果某个源没有产出,目标节点便不会被触发。

FanIn 和上一篇介绍的 Superstep 屏障不是同一个概念。Superstep 屏障负责运行时对并发 executor 的同步,FanIn 屏障则是图结构中指定源节点之间的汇聚条件。

选择边时,可以按照下面的思路判断:

需要分支?
├─ 否 → 直连 AddEdge
└─ 是 → 几路?
    ├─ 2 路,互斥 → 条件边 AddEdge(condition)
    ├─ ≥3 路,选一个 → Switch-Case AddSwitch(配 WithDefault 兜底)
    └─ 一条消息要去多个目标 → FanOut
        ├─ 全部去 → AddFanOutEdge(无 selector)
        └─ 按内容挑几个 → AddFanOutEdge(targetSelector)

需要“等齐了再汇总”?
└─ 是 → FanIn AddFanInBarrierEdge

Executor 的状态与生命周期

上一篇的 Executor 接收输入并返回输出,本身不保存状态。但在实际业务中,节点经常需要维护计数器、缓存或累计结果。下面这个例子来自官方 HITL sample,它会记录用户猜数的次数:

internal sealed class JudgeExecutor(int targetNumber) : Executor<int>("Judge")
{
    private int _targetNumber = targetNumber;
    private int _tries;   // ← 有状态:记试错次数

    public override async ValueTask HandleAsync(int guess, IWorkflowContext context, CancellationToken ct = default)
    {
        this._tries++;
        if (guess == this._targetNumber)
            await context.YieldOutputAsync($"{this._targetNumber} found in {this._tries} tries!");
        else if (guess < this._targetNumber)
            await context.SendMessageAsync(NumberSignal.Below);
        else
            await context.SendMessageAsync(NumberSignal.Above);
    }
}


_tries_targetNumber 都属于节点状态。Workflow 构建后不可变,同一个实例可以重复运行。如果第一次结束时 _tries 等于 5,而节点没有重置,下一次运行就会从 5 继续计数。


这类 Executor 应当实现 IResettableExecutor

public interface IResettableExecutor
{
    ValueTask ResetAsync() => default;   // 默认实现是空操作
}


然后在 ResetAsync 中清空需要跨运行重置的状态:

internal sealed class JudgeExecutor(int targetNumber) : Executor<int>("Judge"), IResettableExecutor
{
    private int _tries;

    public ValueTask ResetAsync()
    {
        this._tries = 0;   // 每次运行前清零
        return default;
    }
    // ... HandleAsync 不变
}


运行结束并释放 workflow 所有权时,框架会调用这些节点的 ResetAsync。如果节点持有可变状态并且会被重复使用,就需要实现这个接口;只有确认节点能够安全地跨运行共享时,才应通过 declareCrossRunShareable: true 明确声明。


除了状态重置,Executor 还需要准确描述自己会发送和产出哪些消息。通过返回值发送消息时,框架可以从方法签名得知类型;如果使用 SendMessageAsyncYieldOutputAsync,尤其是只在某些分支中调用它们,则可以使用 [SendsMessage][YieldsOutput] 显式声明:

[SendsMessage(typeof(NumberSignal))]      // 声明:我会发送 NumberSignal 类型
[YieldsOutput(typeof(string))]            // 声明:我会产出 string 类型的工作流输出
internal sealed class JudgeExecutor() : Executor<int>("Judge") { ... }


声明后,框架可以在构建期把类型注册到节点协议中,对路由和输出进行更准确的校验。

节点内部状态只对当前 Executor 可见。如果多个节点需要共享数据,可以通过 IWorkflowContext 按 scope 读写状态。例如,垃圾检测节点保存邮件原文,下游回复节点再按邮件 ID 读取:

// 写:把邮件存进名为 "EmailState" 的 scope,key 是邮件 ID
await context.QueueStateUpdateAsync(emailId, email, scopeName: "EmailState");

// 读:用同样的 scope + key 取出来
var email = await context.ReadStateAsync<Email>(emailId, scopeName: "EmailState");


scope 相当于状态的命名空间,可以避免不同业务使用相同 key 时发生冲突。共享状态也会进入 checkpoint,因此恢复工作流时可以一并还原。对于体积较大的对象,可以在消息中只传递 ID,需要时再从共享状态中读取。

QueueStateUpdateAsync 不会立即写入,而是在当前超级步结束时应用更新,使状态与超级步边界保持一致。

Executor 还提供了一组生命周期钩子:

钩子触发时机典型用途
InitializeAsyncexecutor 实例创建后、首次执行前初始化资源(连接、缓存)
OnMessageDeliveryStartingAsync每个超级步消息投递准备本步所需资源
OnMessageDeliveryFinishedAsync每个超级步消息投递收尾、刷新缓冲
OnCheckpointingAsync检查点保存前把私有状态写进 checkpoint
OnCheckpointRestoredAsync检查点恢复后从 checkpoint 还原私有状态

对于需要恢复私有状态的节点,OnCheckpointingAsyncOnCheckpointRestoredAsync 尤其重要。例如,HITL 使用的 RequestInfoExecutor 会将尚未处理的外部请求写入 checkpoint,恢复后继续等待响应。


Agent 节点、事件与输出

除了普通 Executor,AIAgent 也可以直接放进工作流。SDK 提供了从 AIAgentExecutorBinding 的隐式转换:

public static implicit operator ExecutorBinding(AIAgent agent) => agent.BindAsExecutor();

因此,将 agent 传给 WorkflowBuilderAddEdge 时,不需要再编写包装类:

AIAgent spamAgent = chatClient.AsAgent(...);
AIAgent replyAgent = chatClient.AsAgent(...);

// spamAgent 直接当起点,replyAgent 直接当目标
var workflow = new WorkflowBuilder(spamAgent)
    .AddEdge(spamAgent, replyAgent)
    .WithOutputFrom(replyAgent)
    .Build();


Agent 节点收到消息后,会调用 agent.RunAsync,并将回复沿边发送。执行过程中除了通用的 ExecutorInvokedEventExecutorCompletedEvent,还会产生两种与 agent 回复有关的事件:

  • AgentResponseEvent 表示一轮完整回复。
  • AgentResponseUpdateEvent 表示流式回复中的一个增量片段,适合用来实时显示生成内容。
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
    if (evt is AgentResponseUpdateEvent update)
    {
        // 增量打印:谁在说 + 这一片的文本
        Console.Write(update.Update.Text);
    }
    else if (evt is AgentResponseEvent response)
    {
        // 完整回复
    }
}

官方 Handoff sample 便是通过监听 AgentResponseUpdateEvent,实时输出当前 agent 及其回复内容。

Agent 节点使用 TurnToken 表示当前轮次可以开始。使用 RunStreamingAsync 并传入普通用户消息时,当前 SDK 会识别工作流是否采用 Chat Protocol,并自动补发一个 TurnToken(emitEvents: true),通常不需要调用方再次发送:

StreamingRun run = await InProcessExecution.RunStreamingAsync(workflow, userInput);


在通过 OpenStreamingAsync 手动控制运行过程,或者实现 GroupChat、Handoff 一类多轮编排时,仍然会直接使用 TurnToken 推进轮次。它的 emitEvents 参数决定是否发布 agent 相关事件。

工作流运行期间产生的所有事件都继承自 WorkflowEvent,主要类型如下:

WorkflowEvent                         (基类,带 Data)
├─ WorkflowStartedEvent               工作流开始
├─ WorkflowOutputEvent                工作流输出(基类)
│  ├─ AgentResponseEvent              agent 完整回复
│  └─ AgentResponseUpdateEvent        agent 流式增量
├─ WorkflowWarningEvent               警告
│  └─ SubworkflowWarningEvent         子工作流警告
├─ WorkflowErrorEvent                 错误
│  └─ SubworkflowErrorEvent           子工作流错误
├─ RequestInfoEvent                   HITL 请求(见后文)
├─ ExecutorEvent                      executor 相关(基类)
│  ├─ ExecutorInvokedEvent            executor 被调用(入参)
│  ├─ ExecutorCompletedEvent          executor 完成(返回值)
│  └─ ExecutorFailedEvent             executor 失败(异常)
└─ SuperStepEvent                     超级步相关(基类)
   ├─ SuperStepStartedEvent           超级步开始
   └─ SuperStepCompletedEvent         超级步完成(含检查点信息!)


调用方可以在 WatchStreamAsync 中按照事件类型分别处理:

await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
    switch (evt)
    {
        case ExecutorInvokedEvent invoked:
            Console.WriteLine($"→ 进入节点 {invoked.ExecutorId},入参: {invoked.Data}");
            break;
        case ExecutorCompletedEvent completed:
            Console.WriteLine($"← 节点 {completed.ExecutorId} 完成,产出: {completed.Data}");
            break;
        case ExecutorFailedEvent failed:
            Console.Error.WriteLine($"✗ 节点 {failed.ExecutorId} 失败: {failed.Data}");
            break;
        case AgentResponseUpdateEvent update:      // agent 流式 token
            Console.Write(update.Update.Text);
            break;
        case WorkflowOutputEvent output:           // 最终输出
            Console.WriteLine($"★ 输出: {output.Data}");
            break;
        case SuperStepCompletedEvent step:          // 超级步完成(带检查点)
            Console.WriteLine($"  (superstep {step.StepNumber} 完成)");
            break;
        case WorkflowErrorEvent err:
            throw err.Data as Exception ?? new InvalidOperationException("workflow 失败");
    }
}


业务代码也可以定义自己的事件。事件类型继承 WorkflowEvent,然后通过 context.AddEventAsync 发布:

internal sealed class DatabaseEvent(string message) : WorkflowEvent(message) { }

// 在 executor 里发
await context.AddEventAsync(new DatabaseEvent($"Email {id} saved to database."));


自定义事件会与内置事件一起从 WatchStreamAsync 返回,调用方可以使用 case DatabaseEvent 单独处理。它不会进入节点之间的消息流,适合记录业务进度和监控信息。

SuperStepCompletedEventCompletionInfo.Checkpoint 还包含当前超级步的检查点,后面恢复流程时会用到。

上一篇提到,WithOutputFrom 决定哪些节点的 yield 会暴露给调用方。完整的输出注册方式如下:

注册方式产出的事件用途
WithOutputFrom(node)不带 tag 的 WorkflowOutputEvent最终结果
WithIntermediateOutputFrom(node)OutputTag.Intermediate 的事件中间进度,例如流式内容或阶段性分析
不注册不产生输出事件仅供工作流内部使用


源码中,_outputExecutors 使用 Dictionary<string, HashSet<OutputTag>> 保存节点 ID 和标签。最终回答节点使用 WithOutputFrom,只用于观察过程的节点使用 WithIntermediateOutputFrom,调用方就能区分最终结果和中间信息。

编排模式 builder(OrchestrationBuilderBase)上也有同名方法,语义一样,只是作用于它内部构建的图。


人工介入与检查点

审批、确认和补充信息等场景需要工作流在中途暂停,等待外部响应后再继续,这就是 Human-in-the-Loop(HITL)。它的核心是 RequestPort:一个可以向调用方请求输入的端口节点。下面仍使用官方的猜数字 sample:

// 定义端口:请求 NumberSignal,响应是 int
RequestPort numberPort = RequestPort.Create<NumberSignal, int>("GuessNumber");
JudgeExecutor judge = new(42);   // 目标数字 42

// 图:端口(起点) → 裁判 → 端口(循环回)
return new WorkflowBuilder(numberPort)
    .AddEdge(numberPort, judge)       // 端口把外部输入送给裁判
    .AddEdge(judge, numberPort)       // 裁判的反馈送回端口(再请求下一次输入)
    .WithOutputFrom(judge)
    .Build();


端口先请求用户猜数,用户输入 int 后由裁判判断。猜错时,裁判向端口发送 NumberSignal.BelowNumberSignal.Above,触发下一轮请求;猜对后则 yield 最终结果。

工作流运行到端口时会暂停并发布 RequestInfoEvent。调用方监听该事件,通过控制台、审批系统或其他外部渠道取得响应,再将结果注入当前运行:

StreamingRun run = await InProcessExecution.OpenStreamingAsync(workflow);

// 1. 启动:发个初始信号,工作流跑到端口就会停
await run.TrySendMessageAsync(NumberSignal.Init);

// 2. 循环:监听事件,遇到请求就问用户、注入答案
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
    if (evt is RequestInfoEvent req)
    {
        // 工作流停下来要输入了。这里简化成问控制台
        Console.Write("猜一个数: ");
        int guess = int.Parse(Console.ReadLine()!);

        // 把响应注入回去,工作流继续
        await run.TrySendMessageAsync(guess);
    }
    else if (evt is WorkflowOutputEvent output)
    {
        Console.WriteLine($"结果: {output.Data}");
        break;   // 猜对了,结束
    }
}


RequestPort.Create<TReq,TResp> 定义请求和响应的类型,StreamingRun.TrySendMessageAsync<T> 则负责从外部注入消息。暂停位置由框架维护,调用方不需要另外记录执行到了哪个节点。

如果等待时间很长,或者进程可能在等待期间退出,还需要通过 checkpoint 持久化执行状态。框架在超级步结束时生成一致的状态快照,检查点的存取由 CheckpointManager 负责:

// 内存版(开发/测试用,进程退出就没了)
CheckpointManager memMgr = CheckpointManager.CreateInMemory();

// JSON 文件版(持久化,生产可用)
using FileSystemJsonCheckpointStore store = new(checkpointFolder);
CheckpointManager jsonMgr = CheckpointManager.CreateJson(store);

// 也有个默认的内存实例
CheckpointManager.Default;


CreateInMemory 适合开发和测试,进程退出后数据便会丢失。CreateJson 接收一个 ICheckpointStore<JsonElement>,例如使用 FileSystemJsonCheckpointStore 将检查点写入磁盘。将 manager 传给执行入口后,运行时会在超级步边界保存检查点:

using FileSystemJsonCheckpointStore store = new(checkpointFolder);
CheckpointManager mgr = CheckpointManager.CreateJson(store);

// 关键:第三个参数传 checkpointManager
StreamingRun run = await InProcessExecution.RunStreamingAsync(workflow, input, mgr);


监听 SuperStepCompletedEvent 可以取得每个超级步对应的检查点信息:

List<CheckpointInfo> checkpoints = [];
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
    if (evt is SuperStepCompletedEvent step && step.CompletionInfo?.Checkpoint is not null)
    {
        checkpoints.Add(step.CompletionInfo.Checkpoint);
        Console.WriteLine($"checkpoint: {step.CompletionInfo.Checkpoint.CheckpointId}");
    }
}


恢复时,将 workflow、检查点和对应的 manager 传给 ResumeStreamingAsync

// 哪怕原进程已经退出,只要 checkpointFolder 里的 JSON 还在,就能从这里继续
StreamingRun resumed = await InProcessExecution.ResumeStreamingAsync(
    workflow, checkpoints[0], mgr);


因此,进程异常退出和人工审批隔天继续都可以使用同一套恢复机制。官方 AotCheckpointing sample 演示了完整过程。

进入 checkpoint 的消息和状态必须能够序列化。共享状态会由框架处理;Executor 的私有状态则需要在 OnCheckpointingAsync 中保存,并在 OnCheckpointRestoredAsync 中恢复。


子工作流与内置编排

当工作流中的节点越来越多时,可以通过 Workflow.BindAsExecutor 将一整个 Workflow 包装成父工作流中的普通节点:

// 子工作流:大写 → 反转 → 加后缀
var subWorkflow = new WorkflowBuilder(uppercase)
    .AddEdge(uppercase, reverse)
    .AddEdge(reverse, appendSuffix)
    .WithOutputFrom(appendSuffix)
    .Build();

// 把子工作流包成一个节点
ExecutorBinding subExecutor = subWorkflow.BindAsExecutor("TextProcessingSubWorkflow");

// 父工作流:加前缀 → (跑子工作流) → 后处理
var mainWorkflow = new WorkflowBuilder(prefix)
    .AddEdge(prefix, subExecutor)
    .AddEdge(subExecutor, postProcess)
    .WithOutputFrom(postProcess)
    .Build();


父工作流只会看到 subExecutor 这个节点。它接收 prefix 的产出,在内部执行完整的子流程,再将输出送给 postProcess。这种方式适合封装和复用相对独立的流程,也能让大型工作流更容易阅读。

同一个子工作流不能同时归属于多个父工作流,成为子工作流后也不能再直接运行。使用 ToMermaidString() 可视化父工作流时,子工作流会显示为单个节点,不会展开内部结构。


顺序、并发、Handoff 和群聊等多 agent 拓扑很常见,SDK 通过 AgentWorkflowBuilder 提供了对应的构建入口:

public static partial class AgentWorkflowBuilder
{
    public static Workflow BuildSequential(params IEnumerable<AIAgent> agents);
    public static Workflow BuildConcurrent(...);
    public static HandoffWorkflowBuilder CreateHandoffBuilderWith(AIAgent initialAgent);
    public static GroupChatWorkflowBuilder CreateGroupChatBuilderWith(Func<...> managerFactory);
    public static SequentialWorkflowBuilder CreateSequentialBuilderWith(params IEnumerable<AIAgent> agents);
    public static ConcurrentWorkflowBuilder CreateConcurrentBuilderWith(params IEnumerable<AIAgent> agents);
    public static MagenticWorkflowBuilder CreateMagenticBuilderWith(AIAgent managerAgent);
}


Sequential 让 agents 依次运行,前一个 agent 的输出会成为后一个的输入,适合翻译、润色、校对这类固定流水线:

Workflow workflow = AgentWorkflowBuilder.BuildSequential(translator, polisher, proofreader);


BuildSequential 还提供带 chainOnlyAgentResponses 参数的重载,用来控制是否只向后传递 agent 回复。


Concurrent 让多个 agent 并行处理同一个输入,适合同一个问题由不同领域的 agent 分别分析:

Workflow workflow = AgentWorkflowBuilder.BuildConcurrent(legalExpert, techExpert, bizExpert);


这些 agent 会在同一个超级步中并行执行,并分别产生输出。


Handoff 从一个入口 agent 开始,根据问题将处理权交给某个专家 agent,适合客服分诊和专家转接:

HandoffWorkflowBuilder builder = AgentWorkflowBuilder.CreateHandoffBuilderWith(intakeAgent);
// 给每个专家配交接关系:从所有其他 agent 都能转到它
foreach (var expert in experts)
    builder.WithHandoffs(allAgents.Except([expert]), expert);

builder.EnableReturnToPrevious();   // 专家处理完能回到上一个 agent
Workflow workflow = builder.Build();


交接通过工具调用触发,框架根据工具调用结果将控制权转给目标 agent。启用 EnableReturnToPrevious 后,专家完成任务还可以回到交接前的 agent。


GroupChat 让多个 agent 围绕共享对话轮流发言,由 GroupChatManager 决定下一位发言者:

// manager 决定发言顺序(比如轮询、或由 LLM 选下一个发言者)
Workflow workflow = AgentWorkflowBuilder
    .CreateGroupChatBuilderWith(managers: agents => new RoundRobinGroupChatManager(agents))
    // 注册参与者、配置...
    .Build();


RoundRobinGroupChatManager 是内置的轮询实现,也可以自定义 manager,通过 LLM 或业务规则选择下一位发言者。


Magentic 属于管理者与执行者模式。manager agent 维护 progress ledger,动态规划并分派任务,适合执行路径无法预先确定的复杂任务。

模式何时用
Sequential固定流水线,前一步喂后一步
Concurrent多视角并行分析同一输入
Handoff分诊/转接,有明确的"前台→专家"关系
GroupChat多方讨论/博弈,需要动态决定发言顺序
Magentic复杂自主任务,需要动态规划任务清单

这些模式最终都会构建成 Workflow,因此仍然可以与自定义 Executor 组合使用。它们只是省去了手工搭建常见拓扑的代码。


查看工作流结构

工作流构建完成后,可以使用 ToMermaidString() 生成 Mermaid 流程图:

Console.WriteLine(workflow.ToMermaidString());

生成结果使用 Mermaid 的 graph 语法,可以直接放进 GitHub、VS Code 或文档站点中查看。条件边和 FanOut 会显示相应标记,子工作流则显示为单个节点。构建完成后先查看一次流程图,通常比运行后再排查错误连线更省时间。

下一篇会把这些能力放进一个完整的 Workflow 示例中,看看条件分流、FanIn、HITL 和检查点如何配合使用。