TradingAgents × OpenClaw — Fork-Join + Yield-Awake 模式深度拆解
1. 整体架构概览
2. 并行 vs 串行 — 两种模式兼有
3. 控制权流转 — sessions_yield 的正确用法
4. 唤醒机制 — 事件驱动的 Promise 模型
5. Wait Table — 唤醒决策的核心数据结构
6. 完整时序图
7. 生活中的类比
8. 应用场景
9. 常见问题 FAQ
10. 一句话总结
多 Agent 系统采用 分层协作 架构。一个 Orchestrator(总调度器) 负责任务拆解与结果汇合,多个专业 子智能体(Sub-agent) 各司其职,底层的 OpenClaw Runtime 负责生命周期管理、唤醒调度。
┌─────────────────────────────────────────────────────────────┐
│ 🧠 Orchestrator(总调度员) │
│ 任务拆解 · 并行派发 · 结果汇合 · 串行决策 │
│ │ │
│ ⬇ 非阻塞 spawn ⬇ │
│ │ │
│ ┌──────────────┬───────┴───────┬──────────────┐ │
│ ▼ ▼ ▼ │ │
│ 📈 基本面 📉 技术面 😊 情绪 │ │
│ 分析师 分析师 分析师 │ │
│ └──────────────┬───────┬───────┘ │ │
│ │ │ │ │
│ ⬇ 全部完成后事件通知 ⬇ │ │
│ ▼ │ │
│ ⚙️ OpenClaw Runtime(Scheduler) │ │
│ 事件驱动 · Wait Table 跟踪 · 全部到齐后唤醒 │ │
│ │ │ │
│ ⬇ 唤醒 ⬇ │ │
│ ▼ │ │
│ 🧠 Orchestrator(恢复执行) │ │
│ │ │ │ │
│ ▼ ▼ │ │
│ ⚠️ 风险管理员 🎯 最终决策 │ │
│ └──────────┘ │ │
│ 🔵 串行执行(有依赖关系) │ │
└─────────────────────────────────────────────────────────────┘
很多人的直觉是"多 Agent = 全部并行",实际上 两种模式都有,取决于任务之间的依赖关系。
🧠 Orchestrator
│
├── 🟢 并行派发(非阻塞 spawn)
│
├── 📈 基本面 📉 技术面 😊 情绪
│
├── 🔵 串行汇合点(yield 等待全部 Promise)
│
├── ⚠️ 风险管理
│
├── 🔵 串行
│
└── 🎯 最终决策
核心原则: 能并行的就并行跑,有依赖关系的必须串行等。这就是 Fork-Join 模型。
原文档中的 await sessions_spawn 会导致串行阻塞,与并行派发的意图冲突。正确的做法是 非阻塞 spawn,将返回的 Promise 收集起来,然后在 yield 时一并交给 Runtime 等待。
T0 · 🟢 Orchestrator 运行中
spawn 三个子智能体(非阻塞,立即返回 Promise)
然后调用 sessions_yield(wait_for=[promise_a, promise_b, promise_c])
│
T1 · 🟡 Orchestrator 挂起
控制权交还给 OpenClaw Runtime
三个子智能体并行运行,Runtime 监控其 Promise 状态
│
T2 · ✅ Agent A 完成
Runtime 收到 A 的完成事件,将对应 Promise resolve
Wait Table 记录 A✅(但 B、C 还没完成 → 继续休眠)
│
T3 · ✅ Agent C 完成
C 的 Promise resolve,Wait Table 更新
(B 仍在跑 → 继续休眠)
│
T4 · ✅ Agent B 完成(最后一个)
B 的 Promise resolve,Wait Table 检测到全部到齐!
Runtime 拿到所有 Promise 的结果,向 Orchestrator 发送
唤醒信号,并附上结果数据
│
T5 · 🟢 Orchestrator 被唤醒
恢复执行,直接获得包含 A、B、C 结果的列表,继续下一步
async def run_analysis():
# 🟢 并行 spawn(非阻塞!没有 await)
promise_a = sessions_spawn(agent="基本面分析师", ...)
promise_b = sessions_spawn(agent="技术面分析师", ...)
promise_c = sessions_spawn(agent="情绪分析师", ...)
# 🟡 让出控制权,同时告知 Runtime:我要等这三个 Promise
# Runtime 负责挂起你,直到所有 Promise 都 resolve
results = await sessions_yield(wait_for=[promise_a, promise_b, promise_c])
# 🟢 被唤醒时,results 已经是 [报告 A, 报告 B, 报告 C]
# 🔵 串行:喂给风险管理员
risk_promise = sessions_spawn(agent="风险管理员", inputs=results)
risk_result = await sessions_yield(wait_for=[risk_promise])
return make_decision(risk_result)
关键变化对比: 从 await sessions_spawn(阻塞)→ 非阻塞 spawn + await sessions_yield(wait_for=[...])(挂起等待)
"是不是 OpenClaw 的进程在实时监控?"——是的。 但这不是"轮询",而是 事件驱动。
while True:
sleep(1 秒) # 空转浪费 CPU
check_all_agents() # 逐个检查
if all_done(): break
# 每个子智能体完成时,Runtime 自动 resolve 对应 Promise
on_promise_resolved(promise, result):
update_wait_table(promise)
if all_promises_resolved(orchestrator):
wake_up(orchestrator, combined_results)
┌────────────────────────────────────────────────────┐
│ ⚙️ OpenClaw Runtime │
│ │
│ 📋 Wait Table 🔗 Promise Registry │
│ 记录谁在等待哪些 存储每个子任务的结果 │
│ │
│ ⏰ Event Loop 🚀 Wake Dispatcher │
│ 异步事件循环 发送唤醒信号+结果 │
│ │
│ 📦 Session Manager │
│ 会话生命周期管理 │
└────────────────────────────────────────────────────┘
子智能体完成任务 → Runtime 将对应 Promise 标记为 resolved,存入结果。
Scheduler 查询 Wait Table → 发现 Orchestrator-A 在等 [Promise-A, Promise-B, Promise-C]。
检查进度:A✅ B❌ C✅ → 不满足,继续等。
稍后 B 完成 → 再次检查:A✅ B✅ C✅ → 全部 resolved!
Wake Dispatcher 向 Orchestrator-A 注入唤醒信号,并 直接带上所有结果。
Orchestrator-A 恢复执行,代码中的 results = await sessions_yield(...) 立即获得数据。
为什么不让 Orchestrator 自己去读 future.result?
因为那样会重复逻辑。既然 Runtime 已经在监控 Promise,就应该在唤醒时直接把结果传回来。这使 yield 的语义更纯粹:挂起我,并给我我需要的结果。
为什么不用"每个 Orchestrator 开一个监控进程"?
因为那是浪费资源。100 个 Orchestrator 就有 100 个空转监控进程。而统一由 Runtime 的 一个 Scheduler 管理所有 Wait Table,实现 多路复用(Multiplexing)——高效、省资源、可扩展。
Wait Table 本质上就是 内存中的一张哈希表,记录了"谁在等哪些 Promise"以及"是否都已 resolved"。
关键改进: 唤醒后直接删除记录而非标记"已唤醒"——避免重复唤醒和内存泄漏。
下面是整个 TradingAgents 多 Agent 交易决策 的完整时序,所有 spawn 均为非阻塞,结果随唤醒一并返回。
Orchestrator Agent-A Agent-B Agent-C Runtime
│ │ │ │ │
├── spawn A────►│ │ │ 注册 Promise-A │
├── spawn B─────│──────────────►│ │ 注册 Promise-B │
├── spawn C─────│───────────────│──────────────►│ 注册 Promise-C │
│ │ │ │ │
│ 🟡 yield(wait_for=[A,B,C]) ────────────────────────────────►│
│ (挂起) │ 🟢 运行 │ 🟢 运行 │ 🟢 运行 │
│ │ │ │ │
│ │ ✅ 完成 │ │ │
│ ├──resolve A───│───────────────│──────────────►│
│ │ │ │ WaitTable: A✅│
│ │ │ │ │
│ │ │ │ ✅ 完成 │
│ │ │ ├──resolve C───►│
│ │ │ │ WaitTable: C✅│
│ │ │ │ │
│ │ │ ✅ 完成 │ │
│ │ ├──resolve B───│──────────────►│
│ │ │ │ ★ 全部 resolved│
│ │ │ ├──唤醒+结果────►│
│◄──── 恢复 + 结果 ───────────────│───────────────│───────────────┤
│ 🟢 拿到3份报告 │ │ │ │
│ │ │ │ │
│ 🔵 串行阶段: │ │ │ │
├── spawn Risk──│───────────────│───────────────│──────────────►│
│ 🟡 yield(wait_for=[Risk]) ────────────────────────────────────►│
│ │ │ 风险管理员运行 │
│ │ │ ✅ 完成 │
│◄──── 唤醒 + 结果 ───────────────│───────────────│───────────────┤
│ 🎯 最终决策 │ │ │ │
▼ ▼ ▼ ▼ ▼
关键改进: Orchestrator 在 yield 时明确告知 Runtime 等待哪些 Promise;唤醒时直接获得结果,无需额外调用 .result。逻辑更简洁,职责更清晰。
关键观察: Orchestrator 在整个过程中经历了 两次 yield → 两次唤醒。第一次等三个并行分析都完成,第二次等风险评估完成。控制权每次都是交还给 Runtime,而不是直接传给下一个 Agent。
老板(Orchestrator)说:"A 去切菜,B 去烧水,C 去摆盘",然后老板 回办公室等着(yield)。
厨房里有个 进度板(Wait Table)。A 切完打个 ✓ → B 烧完打个 ✓ → C 摆完打个 ✓。
最后一个 ✓ 出现 → 传菜员(Runtime)去敲老板的门:"老板,菜、水、盘都好了,可以炒菜了!"
老板出来开始 炒菜(串行阶段)。
主编(Orchestrator)说:"小张去采访现场,小李去查背景资料,小王去约专家评论",然后主编 回办公室喝茶(yield)。
三个人各忙各的。每回来一个人,秘书(Runtime)就在 白板上记一笔。
当白板上三个名字都画了✓ → 秘书提醒主编:"素材都齐了,可以开始写稿了!"
主编开始 写稿(串行)。
指挥(Orchestrator) 说:"小提琴组练旋律,鼓组练节奏,管乐组练和声"
然后指挥 放下指挥棒,坐回椅子上休息(yield)。
三个声部 各练各的。排练助理(Runtime)在门口挂了个 进度表(Wait Table)。
当三个声部全部练完 → 助理通知指挥:"准备齐了,请回来指挥合奏!"
指挥回到台上,开始 指挥合奏(串行阶段)。
多 Agent 的 Fork-Join + Yield-Awake 协同机制适用于广泛的场景。
TradingAgents
一个完整的交易决策需要多个维度的分析:
🧠 Orchestrator
│
├── 🟢 并行 ↓
│
├── 📊 基本面 📉 技术面 😊 情绪面 🌍 宏观
│
├── 🔵 串行汇合 ↓
│
└── ⚠️ 风控 → 🎯 决策 → 📋 执行
效果: 从收到市场事件到生成交易指令,全自动化流水线。每个专业 Agent 独立分析,Orchestrator 汇总后做综合决策。相比单 Agent 串行分析,耗时 减少 60%+(因为三个分析同时跑)。
用户提交一个复杂工单,系统并行派发给多个专业 Agent 处理:
🧠 工单调度器
│
├── 🟢 并行分析 ↓
│
├── 🔍 意图识别 📎 知识库检索 👤 用户画像
│
├── 🔵 串行汇合 ↓
│
└── 💬 回答生成
效果: 用户提问后,意图、知识、用户背景三管齐下同时查,最后生成个性化回复。比传统 chatbot 串行处理快 2-3 倍,且回答质量更高(多个视角综合)。
一篇论文投稿后,需要多个角度的审稿意见:
🧠 编辑调度器
│
├── 🟢 并行审稿 ↓
│
├── 📐 方法学审稿 📊 数据分析审稿 📝 写作与逻辑审稿
│
├── 🔵 串行汇合 ↓
│
└── 📋 综合评审报告
效果: 一篇论文同时由 3 个专业审稿 Agent 评审,编辑只需阅读汇总后的综合报告。审稿周期从 数周缩短到数小时,且覆盖维度更全面。
生产线报警后,需要快速定位故障原因:
🧠 诊断调度器
│
├── 🟢 并行诊断 ↓
│
├── ⚡ 电气系统 Agent 🔧 机械系统 Agent 💻 控制系统 Agent
│
├── 🔵 串行汇合 ↓
│
└── 🚨 故障定位与维修建议
效果: 报警触发后,电气、机械、控制三个子系统同时排查各自领域。全部结果汇合后得出综合诊断结论,平均故障定位时间从 2 小时降至 5 分钟。
面对一个疑似攻击事件,需要多维度交叉验证:
🧠 安全运营中心(SOC)
│
├── 🟢 并行分析 ↓
│
├── 🌐 网络流量 Agent 🖥️ 主机日志 Agent 🔑 身份认证 Agent 🕵️ 威胁情报 Agent
│
├── 🔵 串行汇合 ↓
│
└── 🎯 威胁定级与响应决策
效果: 四个安全 Agent 并行分析同一事件的四个维度。只有 全部维度交叉验证一致 才判定为真实攻击,大幅减少误报。平均响应时间从小时级降至 分钟级。
共同特征: 以上所有场景都符合 Fork-Join 模型——多个独立专业视角同时工作(Fork),到关键决策点汇合(Join),由 Orchestrator 综合各方意见后做出更优质、更可靠的决策。
A: 完全不是一个概念。
A: 是的,Orchestrator 可以在 spawn 时为每个 Promise 设置超时。
# 带超时的 spawn
promise_a = sessions_spawn(agent="基本面分析师", timeout=30) # 30秒后超时
# yield 时也可以设置整体超时
results = await sessions_yield(
wait_for=[promise_a, promise_b, promise_c],
timeout=60 # 60秒后不管是否全部完成都唤醒
)
超时后,未完成的 Promise 会返回一个 TimeoutError 结果,Orchestrator 需要处理这种情况——例如"三个分析中两个已完成,一个超时,是否继续?"
A: 可以!这是 多级 Fork-Join 模式。
🧠 顶层 Orchestrator
│
├── 🧠 子 Orchestrator-A(负责 A 领域)
│ ├── Agent-A1
│ ├── Agent-A2
│ └── Agent-A3
│
├── 🧠 子 Orchestrator-B(负责 B 领域)
│ ├── Agent-B1
│ ├── Agent-B2
│ └── Agent-B3
│
└── 🧠 子 Orchestrator-C(负责 C 领域)
每个子 Orchestrator 内部独立做 Fork-Join,对外暴露 Promise,供父 Orchestrator 等待。架构像一棵树,叶子节点是具体执行任务的子智能体。
A: 纯内存。Wait Table 是 Runtime 进程内部的数据结构,不是持久化存储。原因:
性能要求: 唤醒决策需要在纳秒到微秒级别完成,写数据库太慢
生命周期: Orchestrator 的 yield 是临时的休眠状态,不涉及数据持久化
故障恢复: 如果 Runtime 崩溃,所有等待状态随之消失,这和主流编程语言中的 Future/Promise 设计一致
A: 足够。Wait Table 本质上是一张哈希表,核心操作(插入、更新、检查触发条件)都是 O(1) 时间复杂度。100 个 Orchestrator,每个等 3-5 个 Promise,总条目不超过 500 条,这在现代计算机上几乎没有任何压力。
Fork-Join 模型 = 非阻塞 spawn 派发并行任务 + yield 挂起等待 Promise 全部 resolved + Runtime 事件驱动唤醒并注入结果。
以上分析基于 OpenClaw Runtime 的实际实现,修正版已消除 await sessions_spawn 阻塞、事件/Promise 二义性和 Wait Table 残留状态等问题。