跳到主要内容

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_handleris_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_idtarget_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_idSwitchCaseEdgeGroupCase / ...Default

source/target 去重且保序: source_executor_ids / target_executor_idsdict.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=?)SingleEdgeGroupa→b,可带条件_workflow_builder.py:228
add_fan_out_edges(a, [b,c])FanOutEdgeGroupa 广播给 b、c:280
add_fan_in_edges([a,b], c)FanInEdgeGroupa、b 汇总给 c(c 收 list):509
add_switch_case_edge_group(a, [Case..,Default])SwitchCaseEdgeGroupa 按条件三选一:336
add_multi_selection_edge_group(a, [b,c], selection_func)FanOutEdgeGroupa 按函数选子集:423
add_chain([a,b,c])连串 SingleEdgeGroupa→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_groupsstart_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_iterationasyncio.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.pyExecutorExecutor.executeExecutor._find_handlerExecutor._discover_handlers
handler 装饰器与类型内省_executor.pyhandler_validate_handler_signature
节点类型属性_executor.pyinput_typesoutput_typesworkflow_output_types
边与条件_edge.pyEdgeEdge.should_routeEdgeCondition
边组家族_edge.pyEdgeGroupSingleEdgeGroupFanOutEdgeGroupFanInEdgeGroupSwitchCaseEdgeGroupInternalEdgeGroup
switch-case 条件_edge.pyCaseDefaultSwitchCaseEdgeGroupCaseSwitchCaseEdgeGroupDefault
边上投递_edge_runner.pyEdgeRunnerSingleEdgeRunnerFanOutEdgeRunnerFanInEdgeRunnercreate_edge_runner
fan-in 缓冲/触发_edge_runner.pyFanInEdgeRunner._bufferFanInEdgeRunner._is_ready_to_send
Pregel 超步引擎_runner.pyRunnerRunner.run_until_convergenceRunner._run_iterationRunner._parse_edge_runners
收敛异常 / 上限_runner.py / exceptions.py / _const.pyWorkflowConvergenceExceptionDEFAULT_MAX_ITERATIONS
引擎侧消息/事件总线_runner_context.pyRunnerContextInProcRunnerContextWorkflowMessageMessageType
消息/事件读写_runner_context.pysend_messagedrain_messageshas_messagesdrain_eventsnext_event
节点侧受控接口_workflow_context.pyWorkflowContextWorkflowContext.send_messageWorkflowContext.yield_outputWorkflowContext.add_event
类型推断_workflow_context.pyinfer_output_types_from_ctx_annotationvalidate_workflow_context_annotation
流式构建 API_workflow_builder.pyWorkflowBuilderadd_edgeadd_fan_out_edgesadd_fan_in_edgesadd_switch_case_edge_groupadd_multi_selection_edge_groupadd_chainbuild
输出选择_workflow_builder.py_resolve_designated_executor_idsoutput_fromintermediate_output_from
编译产物 / 运行_workflow.pyWorkflowWorkflow.runWorkflow._execute_with_message_or_checkpointinput_typesoutput_types
图签名_workflow.pyWorkflow._compute_graph_signatureWorkflow._hash_graph_signaturegraph_signature_hash
事件类型_events.pyWorkflowEventsuperstep_startedsuperstep_completedexecutor_invokedexecutor_completedWorkflowEventType