源文件:chapter10/parallel-web-research/README.md 实验 10-6 · 同时从多个网站搜集信息的 Agent(★★) 《深入理解 AI Agent》配套实验。演示多个同构 Agent 的并行搜索 + 中心协调: 主协调器同时启动 N 个子 Agent,每个子 Agent 访问一个"网站/来源"找答案; 一旦某个子 Agent 命中目标,其余立即优雅停止。 书中的原型是"10 个并行的 Computer Use Agent 同时访问不同网站找信息"。为便于 自动验证,本实验不启动真实浏览器,而是用一批可控的模拟信息源代替;把重点 完整放在协调机制上——消息总线、并行派发、实时监控、级联终止、竞态处理均为真实实现。
源文件:chapter10/parallel-web-research/README.md
《深入理解 AI Agent》配套实验。演示多个同构 Agent 的并行搜索 + 中心协调:
主协调器同时启动 N 个子 Agent,每个子 Agent 访问一个"网站/来源"找答案;
一旦某个子 Agent 命中目标,其余立即优雅停止。
书中的原型是"10 个并行的 Computer Use Agent 同时访问不同网站找信息"。为便于
自动验证,本实验不启动真实浏览器,而是用一批可控的模拟信息源代替;把重点
完整放在协调机制上——消息总线、并行派发、实时监控、级联终止、竞态处理均为真实实现。
| 文件 | 作用 |
|---|---|
message_bus.py |
进程内异步消息总线(Redis Pub/Sub 风格),带 Envelope 信封与订阅机制 |
sources.py |
模拟的 10 个"网站/来源",各有不同延迟;可控的关键词命中判断;build_sources(n) 支持任意并行度 |
llm.py |
可选的 LLM 判断层(默认离线关键词判断,配了 key 则用真实大模型) |
agents.py |
主协调器 Coordinator 与子 Agent WorkerAgent(核心协调逻辑);run_sequential 串行基线 |
demo.py |
一条命令的演示入口,带 argparse CLI 与末尾断言式自检 |
┌────────────────────────────┐ │ Coordinator(主协调器) │ │ · 并行派发 task_assigned │ │ · 维护任务状态表(状态机) │ │ · 首个命中→加锁结算(幂等) │ │ · 广播一轮 terminate │ └───────────┬────────────────┘ │ ┌────────────────┴─────── MessageBus(异步消息总线)───────┐ │ Envelope{ sender_id, target, type, payload, seq, ts } │ │ type: task_assigned/status_update/result/terminate/ack │ └──┬─────────┬─────────┬─────────┬──────────────┬──────────┘ │ │ │ │ │ ┌──▼──┐ ┌──▼──┐ ┌──▼──┐ ┌──▼──┐ ... ┌──▼──┐ │W-00 │ │W-01 │ │W-02 │ │W-03 │ │W-09 │ 子 Agent │网站A│ │网站B│ │网站C│ │网站D│ │网站J│ (同构·并行) └─────┘ └─────┘ └─────┘ └─────┘ └─────┘
对应书中强调的五个机制:
typeBUS ... 行就是一次发布/投递(Redis Pub/Sub 语义的进程内实现)。Coordinator 同时给 10 个子 Agent 发 task_assigned 并 asyncio.create_task 并发执行。status_update 上报进度,已提交 → 执行中 →(需要输入)→ 已完成 / 失败 / 已终止。terminate;其余子 Agent 在ack 并优雅退出(状态置为"已终止")。asyncio.Lock +_settled 保证只结算一次、只广播一轮终止;迟到的命中被记录并忽略。为让"竞态""级联终止"可复现,各来源被赋予不同的模拟延迟,其中
geo-journal与forum-qa两个正确源被设成相同延迟,从而稳定地在同一时刻命中、触发竞态。
cd chapter10/parallel-web-research pip install -r requirements.txt # 仅离线演示的话可跳过,纯标准库即可运行 python demo.py
默认走离线关键词判断(无需联网、结果可复现)。若要让子 Agent 用真实 LLM 判断:
cp env.example .env # 在 .env 填入 OPENAI_API_KEY(也支持 Moonshot / 火山方舟 ARK 的 OpenAI 兼容网关) python demo.py # 或不改 .env,直接用命令行开关(仍需配置 key 才会真正生效): python demo.py --use-llm --model gpt-5.6-luna
可用 key:OPENAI_API_KEY(默认模型 gpt-5.6-luna)/ MOONSHOT_API_KEY / ARK_API_KEY
(填到 OPENAI_API_KEY 并按需设置 OPENAI_BASE_URL、OPENAI_MODEL)。
通用回退:若未设置 OPENAI_API_KEY 但设了 OPENROUTER_API_KEY,则真实 LLM 判断
自动改走 OpenRouter,并把模型名映射到其命名空间(gpt-5.6-luna → openai/gpt-5.6-luna)。
不影响协调机制,仅改变"是否命中"的判断。
python demo.py --help 可查看完整帮助。所有参数都不改变默认行为——不传任何参数即为
原有的「10 个 Agent + 内置问题 + 离线可复现 + 详细 BUS 日志」演示。
| 参数 | 作用 | 默认 |
|---|---|---|
-q, --query 问题 |
研究问题(离线关键词判断是针对内置来源调校的,自定义问题一般需搭配 --use-llm) |
内置问题 |
-n, --agents N |
并行子 Agent 数量(N≥2 时始终包含两个含答案的源以稳定演示竞态/级联终止) | 10 |
--model MODEL |
LLM 模型名(等价于设 OPENAI_MODEL,仅 --use-llm 且配置 key 时生效) |
环境变量 |
-o, --output PATH |
把结论(含并行/串行耗时、winner、竞态统计)写入 JSON 文件 | 不写 |
--compare |
并行跑完后再实测一遍串行基线,打印墙钟耗时对比 | 关闭 |
--use-llm |
强制真实 LLM 判断(仍需配 OPENAI_API_KEY 或 OPENROUTER_API_KEY,否则自动回退离线判断) |
关闭 |
--quiet |
减少逐条 BUS 日志(状态表/结论/自检不受影响) | 关闭 |
python demo.py --agents 6 --compare # 6 个并行 Agent,并对比串行墙钟耗时 python demo.py --output result.json # 结论落盘为 JSON
--compare)对应书中实验要求「记录并对比并行/串行时间差异」。--compare 会在并行演示之后,用完全相同的
来源集合再跑一遍串行基线(逐个 await source.fetch() + 判断,命中即止),耗时是实测而非
估算。示例输出(默认 10 源,离线判断):
5) 并行执行墙钟耗时:1.57s(含收敛静默期) ------------------------------------------------------------------------------ 并行 vs 串行 墙钟对比(--compare,串行基线为实测) ------------------------------------------------------------------------------ 串行:命中前逐个抓取了 3/10 个源,墙钟耗时 2.60s,winner=geo-journal 并行:墙钟耗时 1.57s,winner=worker-02 加速比 ≈ 1.66×,节省约 1.03s(并行让最快的源立即结束全局搜索)。
串行必须依次抓完 baike-wiki→news-portal 才轮到最快命中的 geo-journal(累计 2.6s);并行则让
所有源同时开跑,最快的源一命中就触发级联终止、立刻结束全局搜索。并行墙钟包含级联终止的收敛
开销,因此不是理想的「首个源延迟」,但仍显著快于串行——这正是并行 + 级联终止的价值所在。
(a) 消息总线的发布/订阅在工作——每条带信封的消息都打印出来:
BUS [t= 0.00s #3 ] coordinator -> worker-02 | task_assigned | {"question": "...", "source": "geo-journal"} BUS [t= 0.00s #13 ] worker-02 -> coordinator | status_update | {"state": "执行中", ...}
(b) N 个子 Agent 并行执行 + 主 Agent 实时刷新状态表:
── 任务状态表(worker-02 -> 执行中) ── worker-00 源=baike-wiki 状态=执行中 | 开始抓取来源 worker-02 源=geo-journal 状态=执行中 | 开始抓取来源 ...
(c) 级联终止——命中后广播 terminate,其余子 Agent ack 并优雅退出:
BUS [t=0.60s #41 ] coordinator -> ALL | terminate | {"reason":"answer_found","winner":"worker-02"} BUS [t=0.67s #43 ] worker-09 -> coordinator | ack | {"acked":"terminate","source":"map-service"} [ack] worker-09 已确认终止(1 个已 ack) ...最终 8 个未命中的 Worker 全部状态=已终止
(d) 竞态:即使几乎同时命中,也只结算一次、只广播一轮终止:
BUS [t=0.60s #37 ] worker-02 -> coordinator | result | {"found":true, "answer":"...珠穆朗玛峰...8848 米..."} BUS [t=0.60s #38 ] worker-04 -> coordinator | result | {"found":true, "answer":"...珠穆朗玛峰...8848.86 米..."} [结算] 首个命中来自 worker-02 —— 加锁结算,广播一轮 terminate。 [竞态] worker-04 也命中,但已由 worker-02 结算 —— 忽略此次命中,不重复广播终止。
demo.py 末尾有断言式自检:terminate 广播轮数 == 1、只结算一次 == True、winner 非空。跑通即证明机制正确:
4) terminate 广播轮数:1(应为 1,证明只广播一轮) 迟到/并发的重复命中被忽略:['worker-04'] [自检通过] 单次结算 + 单轮终止广播 + 级联 ack 均符合预期。
本实验把重点放在协调机制(消息总线/并行派发/级联终止/竞态处理)上,这些均为真实
实现;但为了可离线运行与自动验证,以下三处做了简化,是已知局限:
WorkerAgent.run() 里的"抓取一步 + 判断"换成真实浏览器geo-journal 与 forum-qa 两个正确源被人为设成MessageBus 用进程内 async 队列模拟 Redis