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-Case | AddSwitch(s, b => b.AddCase(...).WithDefault(...)) | 按顺序匹配多路分支,可以设置默认分支 |
| FanOut | AddFanOutEdge(s, targets, targetSelector) | 将消息发送给一个或多个目标 |
| FanIn | AddFanInBarrierEdge(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 只依赖 input、transform 和 output 三个稳定 ID。不同环境可以为 transform 选择不同实现,也可以从数据库或配置文件读取整组边,再按 ID 注册相应的 Executor。这样改变执行器实现时不需要改动拓扑代码。
需要注意,所有 placeholder 都必须在 Build() 前完成绑定,否则构建时会抛出 InvalidOperationException。Build() 得到的 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 还需要准确描述自己会发送和产出哪些消息。通过返回值发送消息时,框架可以从方法签名得知类型;如果使用 SendMessageAsync 或 YieldOutputAsync,尤其是只在某些分支中调用它们,则可以使用 [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 还提供了一组生命周期钩子:
| 钩子 | 触发时机 | 典型用途 |
|---|---|---|
InitializeAsync | executor 实例创建后、首次执行前 | 初始化资源(连接、缓存) |
OnMessageDeliveryStartingAsync | 每个超级步消息投递前 | 准备本步所需资源 |
OnMessageDeliveryFinishedAsync | 每个超级步消息投递后 | 收尾、刷新缓冲 |
OnCheckpointingAsync | 检查点保存前 | 把私有状态写进 checkpoint |
OnCheckpointRestoredAsync | 检查点恢复后 | 从 checkpoint 还原私有状态 |
对于需要恢复私有状态的节点,OnCheckpointingAsync 和 OnCheckpointRestoredAsync 尤其重要。例如,HITL 使用的 RequestInfoExecutor 会将尚未处理的外部请求写入 checkpoint,恢复后继续等待响应。
Agent 节点、事件与输出
除了普通 Executor,AIAgent 也可以直接放进工作流。SDK 提供了从 AIAgent 到 ExecutorBinding 的隐式转换:
public static implicit operator ExecutorBinding(AIAgent agent) => agent.BindAsExecutor();
因此,将 agent 传给 WorkflowBuilder 或 AddEdge 时,不需要再编写包装类:
AIAgent spamAgent = chatClient.AsAgent(...);
AIAgent replyAgent = chatClient.AsAgent(...);
// spamAgent 直接当起点,replyAgent 直接当目标
var workflow = new WorkflowBuilder(spamAgent)
.AddEdge(spamAgent, replyAgent)
.WithOutputFrom(replyAgent)
.Build();
Agent 节点收到消息后,会调用 agent.RunAsync,并将回复沿边发送。执行过程中除了通用的 ExecutorInvokedEvent 和 ExecutorCompletedEvent,还会产生两种与 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 单独处理。它不会进入节点之间的消息流,适合记录业务进度和监控信息。
SuperStepCompletedEvent 的 CompletionInfo.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.Below 或 NumberSignal.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 和检查点如何配合使用。