Workflow 图引擎:类型路由的 Pregel 超步
30 秒导读: 第 1 章的 Agent 是"让模型边想边做手脚",适合开放式任务;但当你要把好几个 agent、 好几个处理步骤按固定拓扑接起来(先分派、再并行、后汇总),就需要一个确定性的编排底座。 这一章讲的就是这个底座:把编排画成一张有向图——节点是
Executor,边是EdgeGroup, 一个叫Runner的引擎像 Google Pregel 那样一个超步(superstep)一批地把消息从上游推到下游, 直到全图没有消息在流动就算收敛、结束。
本章是全 doc 的技术核心:多 agent 编排(第 5 章)与持久化/人在环路(第 4 章)全都长在这套图引擎之上。 本章只讲图怎么定义、怎么转;检查点、人在环路(HITL)、时间旅行留给 04-durability-hitl.md,高层编排模式(Sequential / Concurrent / GroupChat / Handoff / Magentic) 留给 05-orchestration-patterns.md。Agent 抽象本身见 01-agent-abstraction.md。
1. 这是什么(零基础也能懂)
一句话定义
Workflow 是一张由你提前声明好的有向图:你把每个处理单元(Executor)当节点,用边把它们接起来,
然后引擎负责在节点之间按类型路由消息、分批推进、直到图安静下来。
它和"让模型自由发挥的 agent 循环"是互补关系:
| 维度 | Agent 循环(第 1 章) | Workflow 图(本章) |
|---|---|---|
| 走哪一步 | 模型临场决定 | 图的拓扑提前写死 |
| 确定性 | 低(每次可能不同) | 高(同输入走同路径) |
| 适合 | 开放式、探索性任务 | 固定流程、多 agent 协作 |
| 底层驱动 | tool-call 循环 | 超步批量推进 |
解决什么问题 / 给谁用
假设你要搭一条"客服 工单"流水线:一条消息进来,先分类,高优先级走人工审核 agent、低优先级走自动回复 agent, 两条支线的结果最后汇总成一份报告。这里有分支(switch-case)、并行、汇总(fan-in)—— 用一个大 agent 的 prompt 硬编排既不可靠也没法复用。Workflow 让你把这张图用几行链式代码声明出来, 每个节点关注自己那点逻辑,引擎负责把消息在图里正确地送来送去。
用起来什么样
下面是一条最小工作流:大写化 → 反转,输出结果。感受一下"声明图 → 跑" 的手感。
from typing_extensions import Never
from agent_framework import Executor, WorkflowBuilder, WorkflowContext, handler
class UpperCase(Executor):
@handler
async def process(self, text: str, ctx: WorkflowContext[str]) -> None:
await ctx.send_message(text.upper()) # 发给下游,消息类型 str
class Reverse(Executor):
@handler
async def process(self, text: str, ctx: WorkflowContext[Never, str]) -> None:
await ctx.yield_output(text[::-1]) # 产出为工作流级输出
upper, reverse = UpperCase(id="upper"), Reverse(id="reverse")
workflow = WorkflowBuilder(start_executor=upper).add_edge(upper, reverse).build()
events = await workflow.run("hello")
print(events.get_outputs()) # ['OLLEH']
真实签名见 _workflow_builder.py:60-87(WorkflowBuilder 类 docstring 里的同款示例)。
一句话直觉
把 Workflow 想成一张"消息驿站图"。 每个驿站(Executor)只认某几种"包裹类型",
来了包裹就按类型挑对应的处理程序拆开处理,处理完可能生成新包裹丢回邮路;驿站之间的路(Edge)决定包裹往哪送;
邮政系统(Runner)不是来一个送一个,而是攒一批、统一发一轮(一个超步),一轮轮发到邮路上再没有包裹为止。
本节不碰底层。记住三个词就够往下走:节点按类型路由、边定拓扑、runner 按超步推进。
2. 顶层全景(它大概怎么转)
整个引擎分两段生命周期:先由 WorkflowBuilder 把图编译成不可变的 Workflow,再由 Runner 在运行时驱动它。
部件一句话职责
| 部件 | 干什么 | 在哪个文件 |
|---|---|---|
Executor | 图的节点;按消息类型把消息路由到对应 @handler | _executor.py |
@handler | 装饰器;标记某方法处理某类型消息,并记录输入/输出类型 | _executor.py:547 |
Edge / EdgeGroup | 边与边组;定义"谁连谁"以及 fan-out/fan-in/switch-case 等路由语义 | _edge.py |
EdgeRunner | 边上的投递器;做类型/条件过滤后把消息交给目标节点执行 | _edge_runner.py |
WorkflowBuilder | 流式 API;add_edge 等一路声明,build() 编译并校验成 Workflow | _workflow_builder.py |
Workflow | 编译产物(不可变);持有 executors + edge_groups + 图签名,暴露 run | _workflow.py |
Runner | "Pregel 超步"引擎;run_until_convergence 一批批推进消息直到收敛 | _runner.py:45 |
RunnerContext | 运行时消息/事件总线;send_message / drain_messages / drain_events | _runner_context.py |
WorkflowContext | 递给 handler 的受控接口;send_message / yield_output / add_event | _workflow_context.py |
WorkflowEvent | 统一事件类型;超步、executor 生命周期、输出全走它 | _events.py:146 |
图 A:编译期(声明 → 不可变产物)
从左到右是一次 build()。要点:每加一个 executor,builder 会自动补一条 InternalEdgeGroup(内部消息用);
build() 里先跑图校验,再把运行时对象组装进 Workflow。
WorkflowBuilder(流式声明) build() Workflow(不可变编译产物)
───────────────────────── ─────────────── ──────────────────────────
start_executor=upper executors: {id -> Executor}
.add_edge(upper, reverse) ──► validate_workflow_graph edge_groups: [EdgeGroup...]
.add_fan_out_edges(...) (类型/连通性校验) graph_signature + sha256 hash
.add_switch_case_edge_group(...) │ ┌───────────────┐
output_from=[...] └── 通过 ──► │ Runner(内部) │
└───────────────┘
图 B:运行期(一个超步怎么走)
这是本章的心脏。怎么读: 一个超步 = "抽干上一轮所有在途消息 → 按来源分批并发投递 → 目标节点跑 handler、 可能产出新消息 → 提交状态 → 还有消息就进下一超步"。
超步 N
│
├─ ctx.drain_messages() ──► { upper: [msg], router: [m1,m2] } ← 按「发送方 id」分桶
│ │
│ │ 每个桶并发;桶内每个 EdgeRunner 保序
│ ▼
├─ EdgeRunner.send_message(msg)
│ │ ① 目标能处理这类型吗?(can_handle) ② 边条件为真吗?(should_route)
│ │ 都过 ──► executor.execute(msg)
│ ▼ │ handler 内:
│ 目标节点跑 handler ├─ ctx.send_message(x) ─► 入队,等下一超步再投
│ └─ ctx.yield_output(y) ─► 立即变 output 事件
│
├─ state.commit() + (若开启)create_checkpoint
│
└─ 还有在途消息?── 是 ──► 超步 N+1
└─ 否 ──► 收敛,run 结束
一句话主线:输入喂给 start executor → 它发消息 → runner 一超步把消息投给下游 → 下游再发 → …… → 没消息可投 = 收敛。
3. 核心原理(逐个机制,由浅入深)
3.1 Executor:按消息类型路由的节点
它要解决的小问题: 一个节点可能要处理好几种消息(str、某个 dataclass、某个 response)。
怎么让"来了什么类型就自动跑对应逻辑",而不用手写一堆 if isinstance?
思路: 用装饰器 @handler 标注每个处理方法,从它的类型注解里抽出"这个方法吃什么类型"。
节点初始化时扫描自己所有带标记的方法,建一张 {消息类型 -> 方法} 的表;来消息时按类型查表分发。
原理演示(示意,非源码):
# 演示「按类型建表 + 查表分发」这一个想法
class Node:
def __init__(self):
self.handlers = {} # 类型 -> 方法
for name in dir(self):
fn = getattr(self, name)
spec = getattr(fn, "_handler_spec", None)
if spec: # 被 @handler 标记过
self.handlers[spec["message_type"]] = fn
async def execute(self, msg):
for t, fn in self.handlers.items():
if isinstance(msg, t): # 第一个类型匹配的 handler 胜出
return await fn(msg)
raise RuntimeError("no handler for this type")
真实实现:
- 建表:
_discover_handlers遍历dir(self.__class__),凡带_handler_spec属性的可调用就登记进self._handlers, 重复类型直接报错(_executor.py:329-353)。构造时若没发现任何 handler 也报错(_executor.py:210-214)。 - 装饰器:
handler支持两种取类型的模式——内省(从注解推,默认)与显式(@handler(input=..., output=...), 一旦给了任一显式参数就完全关闭内省)。它把结果塞进wrapper._handler_spec(_executor.py:547-677)。 - 分发:
_find_handler用is_instance_of(message, message_type)逐个匹配,返回第一个命中的 handler; 匹配不到抛RuntimeError(_executor.py:455-491)。 - 执行:
execute是引擎调用、你不要override的入口——它开 tracing span、找 handler、必要时把WorkflowMessage拆包成裸数据、造WorkflowContext,发executor_invoked事件、await handler(...)、再发executor_completed(_executor.py:219-294)。
类型即契约: handler 的 WorkflowContext[OutT, W_OutT] 注解不只是给 IDE 看的——引擎从中推断这个节点的
输出消息类型(OutT)和工作流级输出类型(W_OutT),暴露成 input_types / output_types /
workflow_output_types 三个属性(_executor.py:410-449),build() 时用来做类型兼容校验。推断逻辑在
infer_output_types_from_ctx_annotation(_workflow_context.py:35-112)。
3.2 Edge 与 EdgeGroup:边定义拓扑
它要解决的小问题: 节点有了,怎么表达"A 之后走 B""A 广播给 B、C""B、C 汇总到 D""按条件三选一"?
思路: 最小单位是 Edge(一条有向边 + 可选布尔条件)。但引擎不直接摆弄裸边,而是把语义相同的边打包成
EdgeGroup,让 runner 能对"fan-out / fan-in / switch-case"这些高阶路由整体推理。
一条 Edge 长什么样: 只存 source_id、target_id、可选 condition;id 是 "source->target";
should_route(data) 跑条件(无条件恒为 True,支持 async)(_edge.py:75-198)。注意条件收的是消息数据本身,
返回该不该走这条边。
六种 EdgeGroup:
| 边组 | 语义 | 约束 | 类 |
|---|---|---|---|
SingleEdgeGroup | 一对一(可带条件) | — | _edge.py:468 |
FanOutEdgeGroup | 一对多广播,可选 selection_func 缩小目标 | 至少 2 个目标 | _edge.py:499 |
FanInEdgeGroup | 多对一汇聚,目标收到 list | 至少 2 个源 | _edge.py:614 |
SwitchCaseEdgeGroup | 多分支,按条件顺序命中一支 | 恰好 1 个 Default | _edge.py:806 |
InternalEdgeGroup | 系统内部消息(request/response 用),无条件 | 每个 executor 自动配一条 | _edge.py:908 |
| (multi-selection) | 复用 FanOutEdgeGroup + 自定义 selection_func | 见 3.5 | _edge.py:499 |
switch-case 的巧处: 它继承自 FanOutEdgeGroup——把"分支选择"实现成一个特殊的 selection_func:
按顺序试每个 Case.condition,第一个为真的返回其目标;遇到 Default 直接返回其目标;都不中就抛错
(_edge.py:862-871)。所以 switch-case 本质是"目标只选一个"的 fan-out。运行时用 Case / Default
两个轻量壳承载"条件 + 目标 executor"(_edge.py:244-291),持久化时才转成只存 target_id 的
SwitchCaseEdgeGroupCase / ...Default。
source/target 去重且保序: source_executor_ids / target_executor_ids 用 dict.fromkeys 保留首次出现顺序
(_edge.py:351-376)——runner 建"发送方 → 边组"映射时依赖这个确定性顺序。
可序列化: 每个 EdgeGroup 都能 to_dict / from_dict,并用 _TYPE_REGISTRY 按类名恢复子类
(_edge.py:398-465)。回调(条件、selection_func)不持久化——反序列化时装一个"一调用就响亮报错"的占位
_missing_callable(_edge.py:50-72),逼你重新注册,而不是静默失效。
3.3 EdgeRunner:边上的消息投递
它要解决的小问题: 有了边组的"静态语义",谁在运行时真正把一条消息过滤 + 投递到目标? fan-in 还得攒齐多个源才能触发,这状态存哪?
思路: 每种边组配一个 EdgeRunner(create_edge_runner 工厂按类型选,_edge_runner.py:411-429)。
投递前统一做两道关:目标能处理这类型吗(_can_handle)、边条件为真吗(should_route);都过才
_execute_on_target 真正调 executor.execute(_edge_runner.py:60-88)。每次投递还记 OpenTelemetry span,
把"投没投、为什么没投"(类型不匹配 / 条件为假 / 目标不符)写进 EdgeGroupDeliveryStatus。
四种 runner 的关键差异:
SingleEdgeRunner:一条边,过两关就投(_edge_runner.py:91-160)。FanOutEdgeRunner:先跑selection_func选出目标子集,再对每个目标各过两关;多个目标用asyncio.gather并发投递(_edge_runner.py:279-287)。若消息带了具体target_id,只投那一个。SwitchCaseEdgeRunner:直接继承FanOutEdgeRunner,靠边组里那个"按 case 选一支"的 selection_func 干活 (_edge_runner.py:404-408)。FanInEdgeRunner:有状态。按源 id 把消息缓冲进self._buffer;每次收消息后查_is_ready_to_send(所有源都到齐了吗);齐了才把各源数据聚合成一个 list,造一条新的聚合WorkflowMessage投给目标,再清空缓冲 (_edge_runner.py:297-401)。没齐就先 buffer 着、返回True(表示"已受理,等更多")。
一个易踩的点: fan-in 判断目标"能不能处理"时,是拿 [message.data](包成 list)去问 can_handle
(_edge_runner.py:332-334)——因为聚合后目标收到的是 list,所以目标 handler 的输入类型必须是 list[...]。
3.4 Runner:Pregel 超步,推进到收敛
它要解决的小问题: 节点会在处理中发新消息,新消息又触发新节点……怎么有序推进这个"消息生消息"的过程, 既不乱序、又能并行、还能判断"什么时候算跑完了"?
思路:借 Google Pregel 的超步模型。 不是"来一条投一条",而是分轮:每一轮(超步)只处理"上一轮结束时 已经在途的所有消息";这一轮里产生的新消息攒着,留到下一轮统一投。一轮结束提交状态、存检查点; 某一轮开始时发现没有在途消息了,就叫收敛,结束。
主循环骨架(示意,非源码,浓缩自 run_until_convergence):
# 演示「超步推进 + 收敛判据」这一个想法
while self._iteration < self._max_iterations:
yield superstep_started(self._iteration + 1)
await self._run_iteration() # 抽干这一轮消息、并发投递、产出新消息入队
self._iteration += 1
self._state.commit() # 超步边界统一提交状态
await self.create_checkpoint_if_enabled()
yield superstep_completed(self._iteration)
if not await self._ctx.has_messages(): # ← 收敛判据:没有在途消息
break
# 超了上限还有消息 → 判为不收敛
if self._iteration >= self._max_iterations and await self._ctx.has_messages():
raise WorkflowConvergenceException(...)
真实实现: run_until_convergence(_runner.py:107-184)。几个关键细节:
- 超步 0 特殊处理: start executor 是在循环外先跑的(见 3.7),所以循环外若已有消息且是第 0 轮,
先存一个"超步 0 结束"的检查点(
_runner.py:120-121)。 - 边跑边流事件: 迭代协程用
asyncio.create_task起,主循环一边next_event()轮询、一边让迭代推进, 实现"事件实时流出"而非攒到超步末尾(_runner.py:129-144)。 - 收敛上限:
max_iterations默认 100(DEFAULT_MAX_ITERATIONS,_const.py:4);超了还有消息就抛WorkflowConvergenceException(exceptions.py:251)——这是防死循环的护栏。
一个超步内部(_run_iteration,_runner.py:186-236)的并发/保序规则:
drain_messages() ──► { 发送方A: [m1,m2], 发送方B: [m3] }
│
├─ 不同「发送方」的桶 ── 并发(asyncio.gather)
│
└─ 同一发送方 → 它的多个 EdgeRunner ── 并发
└─ 同一 EdgeRunner 内多条消息 ── 严格保序(for 循环顺序投)
即:跨边组并行,单边组内保序。这保证了"同一条路径上的消息顺序"不乱,同时榨取了并行度。
_edge_runner_map(发送方 id → 该发送方的所有 EdgeRunner)在构造时由 _parse_edge_runners 建好(_runner.py:383-398)。
3.5 消息与事件的两条总线:RunnerContext 与 WorkflowContext
它要解决的小问题: handler 里 ctx.send_message(x) 之后,消息去哪了?它凭什么下一超步能被投出去?
事件(输出、日志)又怎么流到调用方?
思路:分两层 context。
RunnerContext(引擎侧总线,_runner_context.py:97):真正存消息和事件的地方。默认实现InProcRunnerContext把消息存进一个{发送方 id -> [WorkflowMessage]}的字典,事件存进一个asyncio.Queue。WorkflowContext(节点侧受控接口,_workflow_context.py:207):递给每个 handler 的ctx,是对 RunnerContext 的收窄封装——只暴露send_message/yield_output/add_event/request_info/ 读写 state,不给节点碰引擎内部。
一条消息的旅程:
handler 内 ctx.send_message(x) [WorkflowContext.send_message, _workflow_context.py:308]
│ 包成 WorkflowMessage(data=x, source_id=本节点id)
▼
RunnerContext.send_message(msg) [InProcRunnerContext, _runner_context.py:303]
│ 存进 _messages[本节点id].append(msg) ← 攒着,本超步不投
▼
下一超步 drain_messages() 抽干这批 [_runner_context.py:307]
│
▼
按 source_id 找到对应 EdgeRunner,过滤后投给下游 [Runner._run_iteration]
消息 vs 事件的分工:
- 消息(message) 走
_messages字典,是节点间的数据流,不出现在事件流里——调用方看不到中间消息。 - 事件(event) 走
_event_queue,是流给调用方的东西:超步开始/结束、executor 生命周期、yield_output的输出、 warning/error。add_event/drain_events/next_event(_runner_context.py:315-341)。
WorkflowMessage 本身(_runner_context.py:36-94)带 source_id / target_id / type(STANDARD 或 RESPONSE)
以及 OpenTelemetry trace 上下文(用 list 支持 fan-in 多源聚合)。
yield_output 的门道: 它不只是"发个输出事件",还要按工作流的输出选择策略给事件贴标签——
output-designated 节点产 type='output',intermediate-designated 产 type='intermediate',没选中的直接隐藏
(_workflow_context.py:340-370,分类器是 classify_yielded_output)。这套 output/intermediate 选择由 builder 的
output_from / intermediate_output_from 参数决定(见 3.6)。
add_event 的护栏: 节点想直接发 output/intermediate/started/status/failed 这些保留类型的事件,
会被拦下并转成 warning(_workflow_context.py:372-391)——输出只能走 yield_output,生命周期事件只能框架发。
3.6 WorkflowBuilder:流式声明图
它要解决的小问题: 怎么让"声明这张图"读起来像说人话,同时在 build() 时就把类型不兼容、图不连通这些错误挡掉?
思路: 一组返回 self 的链式方法,每个方法往内部 _edge_groups 里塞一个对应的 EdgeGroup;build() 收尾做校验、
造运行时对象、封成不可变 Workflow。
流式 API 一览:
| 方法 | 塞进的边组 | 语义 | 位置 |
|---|---|---|---|
add_edge(a, b, condition=?) | SingleEdgeGroup | a→b,可带条件 | _workflow_builder.py:228 |
add_fan_out_edges(a, [b,c]) | FanOutEdgeGroup | a 广播给 b、c | :280 |
add_fan_in_edges([a,b], c) | FanInEdgeGroup | a、b 汇总给 c(c 收 list) | :509 |
add_switch_case_edge_group(a, [Case..,Default]) | SwitchCaseEdgeGroup | a 按条件三选一 | :336 |
add_multi_selection_edge_group(a, [b,c], selection_func) | FanOutEdgeGroup | a 按函数选 子集 | :423 |
add_chain([a,b,c]) | 连串 SingleEdgeGroup | a→b→c 顺次接 | :564 |
build() | — | 校验 + 编译成 Workflow | :725 |
几个不显然的设计:
- agent 自动包装: 方法参数既收
Executor也收SupportsAgentRun(一个 agent);后者会被_maybe_wrap_agent自动包成AgentExecutor,且同一个 agent 实例复用同一个 wrapper(_workflow_builder.py:189-226), 避免同一 agent 被包成多个节点。这就是"agent 无缝当节点用"的入口(细节见第 5 章)。 - 每个节点自动配内部边组:
_add_executor每登记一个新节点,就补一条InternalEdgeGroup(:172-187)—— 这条边专门送 request/response 这类系统内部消息(HITL 用,见第 4 章)。 - 输出选择: 构造器的
output_from/intermediate_output_from决定哪些节点的yield_output算"最终输出"、 哪些算"中间输出"、其余隐藏。两者都不给会走兼容模式(每个 yield 都当 output)并发弃用警告(:815-823);"all"/"all_other"/ 显式 list 三种粒度,规则在 docstring:128-140和解析函数_resolve_designated_executor_ids(:660-682)。 - 校验前置:
build()里validate_workflow_graph检查起点已设、边连的是有效节点、图连通、相邻节点类型兼容 (:849-856),不通过直接抛WorkflowValidationError,把错误挡在运行前。
3.7 Workflow:编译产物,怎么跑、怎么锁拓扑
它要解决的小问题: build() 出来的 Workflow 里有什么?run 怎么把输入喂进图?checkpoint 恢复时怎么保证
"图没变过"?
编译产物里有: executors(id→节点)、edge_groups、start_executor_id、以及一个图签名及其 sha256 哈希
(_workflow.py:323-335)。内部还揣着一个 Runner 和一个 RunnerContext(:349-359)。Workflow 是不可变的、
可多次 run(状态在实例内跨 run 保留,要独立运行就另建实例)。
输入怎么进图: run 最终走到 _execute_with_message_or_checkpoint——新消息时直接调
start executor 的 execute,把输入喂进去(_workflow.py:657-666)。注意这一步在超步循环之外发生:
start executor 先跑一遍、产出第一批消息,循环再从超步 1 开始把它们推下去。
输入/输出类型: 工作流的 input_types 就是 start executor 的 输入类型(:1117-1127);
output_types 是所有节点 workflow_output_types 的并集(:1129-1145)。
图签名 = 拓扑指纹: _compute_graph_signature 把"起点 + 每个节点的类全名 + 每条边组的类型/源/目标/边/选择函数名"
规范化成一个 dict,排序后 json.dumps 再 sha256(_workflow.py:1050-1115)。它只含结构、不含数据/状态。
用途:从 checkpoint 恢复时比对哈希——拓扑变了就拒绝恢复(_runner.py:316-320),防止"用旧快照喂给改过的图"。
子工作流(WorkflowExecutor)的签名会递归嵌进父签名(:1064-1068),所以内层图改了外层也能察觉。
4. 深入实现(串一条端到端路径)
把 3.x 拼起来,追一条真实路径:一条 str 进入"分派 → 并行两 agent → 汇总"的图。
run("工单内容")
│
├─[循环外] start=Dispatcher.execute("工单内容") _workflow.py:657
│ handler 里 ctx.send_message(Task(...)) _workflow_context.py:308
│ └─ 存进 _messages["dispatcher"] _runner_context.py:303
│
├─ 超步 1: _run_iteration _runner.py:186
│ drain_messages() → {"dispatcher":[Task]} _runner_context.py:307
│ dispatcher 的 FanOutEdgeRunner.send_message(Task) _edge_runner.py:175
│ selection_func 选中 [agentA, agentB]
│ 两目标各过 can_handle + should_route,asyncio.gather 并发投 _edge_runner.py:279
│ agentA.execute(Task) / agentB.execute(Task) _executor.py:219
│ 各自 ctx.send_message(Partial(...)) → 入队
│ state.commit() (+ checkpoint) _runner.py:165
│ 还有消息 → 进超步 2
│
├─ 超步 2:
│ drain → {"agentA":[Partial], "agentB":[Partial]}
│ 两桶并发;各自的 FanInEdgeRunner 把消息 buffer 进 _buffer _edge_runner.py:336
│ agentA 那条:_is_ready_to_send? 还差 agentB → 先 buffer
│ agentB 那条:齐了 → 聚合成 [PartialA, PartialB] 投给 Aggregator _edge_runner.py:349
│ Aggregator.execute([Partial,Partial])
│ ctx.yield_output(报告) → 立即成 type='output' 事件 _workflow_context.py:340
│ 没有在途消息了
│
└─ 超步 3 开始前:has_messages()==False → 收敛,run 结束 _runner.py:173
get_outputs() 收集所有 output 事件的 data
这条路径把三件事演全了:类型路由(每个 execute 内部按类型找 handler)、边定拓扑(fan-out 分发、fan-in 汇聚)、 超步推进(消息攒到超步边界统一投、fan-in 跨超步攒齐才触发、无消息即收敛)。
5. 巧妙之处(可借鉴的技术)
-
switch-case 复用 fan-out。 不为分支单独造一套投递机制,而是把"选哪一支"实现成一个特殊的
selection_func塞进 fan-out——SwitchCaseEdgeGroup(FanOutEdgeGroup)+SwitchCaseEdgeRunner(FanOutEdgeRunner)(_edge.py:806、_edge_runner.py:404)。一套投递代码覆盖"广播 / 子集 / 单选"三种语义。 -
类型注解当运行时契约。 handler 的
WorkflowContext[OutT, W_OutT]泛型参数被引擎读出来做类型校验和输出推断 (_workflow_context.py:35-112),不是纯静态摆设。声明即约束。 -
超步边界 = 天然的一致性/检查点点位。 消息只在超步之间投、状态只在超步末
commit、检查点也在超步末存 (_runner.py:164-168)。"批量推进"这一个决策同时买到了确定性、可并行、可持久化三样(持久化细节见第 4 章)。 -
回调不持久化,占位符响亮失败。 序列化图时条件/选择函数换成
_missing_callable(_edge.py:50-72), 一被调用就抛"序列化后不可用",逼你显式重注册,而不是静默走错分支。 -
图签名递归嵌套。 子工作流签名嵌进父签名(
_workflow.py:1064),让"内层拓扑变化"也能在外层 checkpoint 恢复时被拦下。 -
跨边组并行、单边组内保序。
_run_iteration用asyncio.gather榨并行,又用"同一 EdgeRunner 内 for 顺序投"保住 单路径顺序(_runner.py:186-236),在并发与顺序保证间取了个务实的平衡点。
6. 边界与局限
-
不是真并行。 全程
asyncio单线程事件循环;"并发投递"是协程交错,不是多核并行——源码注释明说"true parallelism is not realized in Python"(_runner.py:197)。CPU 密集节点仍会阻塞。 -
收敛靠
max_iterations兜底。 环形拓扑或持续自发消息会一直不收敛,只能靠默认 100 的上限抛WorkflowConvergenceException(_runner.py:178)。上限是护栏,不是"检测到无限循环"。 -
fan-in 要全员到齐。
_is_ready_to_send要求每个源都有 buffer 消息才触发(_edge_runner.py:399-401); 某个源这一轮没产消息,汇总就一直等——没有"超时/部分聚合"的内建语义。 -
同类型 handler 唯一。 一个节点里两个 handler 吃同一类型会在建表时报
Duplicate handler(_executor.py:346); 分派是"第一个类型匹配即用",子类型/联合类型重叠时要留意匹配顺序。 -
中间消息对调用方不可见。 只有 output/intermediate/生命周期/状态事件进事件流;节点间消息不进(
_workflow.py:223-225)。 想观测中间数据得靠intermediate_output_from显式暴露,或 tracing span。 -
Runner是内部 API。 虽可从agent_framework导入,但已标记弃用、仅供内部用(_runner.py:30-42); 正常路径永远经Workflow.run。
7. 横向对比
同 shelf 的 VoltAgent 也有 workflow 引擎(见 voltagent/05-workflow-engine.md),但取舍不同:
| 维度 | Microsoft Agent Framework(本章) | VoltAgent Workflow |
|---|---|---|
| 结构模型 | 有向图 + 超步(Pregel) | 线性链 + 步骤原语(andThen/andAgent/andAll…) |
| 推进方式 | 消息驱动、批量超步、按类型路由 | 步骤顺序执行,数据在步间传递 |
| 并发 | fan-out / fan-in 边组 + 超步内并发 | andAll / andRace 并行原语 |
| 分支 | switch-case 边组(复用 fan-out) | 步骤内条件逻辑 |
| 停/恢复 | 超步末检查点 + 图签名校验(第 4 章) | 每步存快照,可暂停/跨进程恢复/重放 |
一句话:MAF 用"图 + 超步"换来了 fan-in 汇聚与多 agent 拓扑的天然表达;VoltAgent 用"链 + 原语"换来了更线性直观的 声明手感。 需要复杂并行/汇聚拓扑选前者,需要顺序流水线选后者。
8. 代码地图(导航索引)
| 主题 | 文件 | 关键符号 |
|---|---|---|
| 节点基类 / 执行入口 | _executor.py | Executor、Executor.execute、Executor._find_handler、Executor._discover_handlers |
| handler 装饰器与类型内省 | _executor.py | handler、_validate_handler_signature |
| 节点类型属性 | _executor.py | input_types、output_types、workflow_output_types |
| 边与条件 | _edge.py | Edge、Edge.should_route、EdgeCondition |
| 边组家族 | _edge.py | EdgeGroup、SingleEdgeGroup、FanOutEdgeGroup、FanInEdgeGroup、SwitchCaseEdgeGroup、InternalEdgeGroup |
| switch-case 条件 | _edge.py | Case、Default、SwitchCaseEdgeGroupCase、SwitchCaseEdgeGroupDefault |
| 边上投递 | _edge_runner.py | EdgeRunner、SingleEdgeRunner、FanOutEdgeRunner、FanInEdgeRunner、create_edge_runner |
| fan-in 缓冲/触发 | _edge_runner.py | FanInEdgeRunner._buffer、FanInEdgeRunner._is_ready_to_send |
| Pregel 超步引擎 | _runner.py | Runner、Runner.run_until_convergence、Runner._run_iteration、Runner._parse_edge_runners |
| 收敛异常 / 上限 | _runner.py / exceptions.py / _const.py | WorkflowConvergenceException、DEFAULT_MAX_ITERATIONS |
| 引擎侧消息/事件总线 | _runner_context.py | RunnerContext、InProcRunnerContext、WorkflowMessage、MessageType |
| 消息/事件读写 | _runner_context.py | send_message、drain_messages、has_messages、drain_events、next_event |
| 节点侧受控接口 | _workflow_context.py | WorkflowContext、WorkflowContext.send_message、WorkflowContext.yield_output、WorkflowContext.add_event |
| 类型推断 | _workflow_context.py | infer_output_types_from_ctx_annotation、validate_workflow_context_annotation |
| 流式构建 API | _workflow_builder.py | WorkflowBuilder、add_edge、add_fan_out_edges、add_fan_in_edges、add_switch_case_edge_group、add_multi_selection_edge_group、add_chain、build |
| 输出选择 | _workflow_builder.py | _resolve_designated_executor_ids、output_from、intermediate_output_from |
| 编译产物 / 运行 | _workflow.py | Workflow、Workflow.run、Workflow._execute_with_message_or_checkpoint、input_types、output_types |
| 图签名 | _workflow.py | Workflow._compute_graph_signature、Workflow._hash_graph_signature、graph_signature_hash |
| 事件类型 | _events.py | WorkflowEvent、superstep_started、superstep_completed、executor_invoked、executor_completed、WorkflowEventType |