M2 · 同质 Fork 与并发控制#
简历 Bullet Point: 实现同质子 Agent fork 并行调度(主/子共享完整工具集与 system prompt,闭包延迟求值解循环依赖),三层并发控制(fork 次数预算 + Semaphore 并发度 + 树级检索预算),ContextVar 做隔离、module-level dict 做跨 Task 聚合;安全四层确保弱模型不失控
开场钩子#
场景#
第一版跨平台并行 fork 做出来就炸了。子 Agent 握着完整工具集——包括 dispatch_tool 自身,一个子觉得结果不够好又 fork 孙,资源指数爆炸。更隐蔽的是检索预算用 ContextVar 做计数器,5 个子各搜 3 次,父 loop 计数器还是 0——asyncio 创建子任务拷贝 context,子 set 的值父读不到。
逼出两条核心设计:安全四层先于功能 + 隔离用 ContextVar、聚合用 module dict。
面试官切入#
“你说’同质 fork’——子和父到底什么相同什么不同?全给一样的工具不怕子 Agent 乱调吗?“
一、模块运作流程#
1.1 一句话定位#
同质 fork 让主 Agent 按需复制自己并行执行子任务,只回传精简结论;三层并发控制 + 安全四层保证弱模型不无限繁殖。
1.2 全景流程图#
主 AgentLoop (depth=0)
├─ Think: 判断满足 fork 三件事之"能并行"
├─ Act: parallel_dispatch_tool(demands_list=[5 条平台子目标])
│ ├─ _ensure_platform_coverage(): 模型只写 3 条 → 机制补齐到 5 条
│ ├─ ForkBudget.charge() → 放行
│ ├─ asyncio.gather(5 × _run_sub_agent, return_exceptions=True)
│ │ ┌──────────────────────────────────────────┐
│ │ │ 子 Agent #1 (depth=1) │
│ │ │ ├─ enter_fork() → depth 0→1 ①深度闸 │
│ │ │ ├─ wait_for(timeout=90s) ②超时 │
│ │ │ ├─ item_search(platform="amazon") │
│ │ │ │ truncate(4000 tok) ③截断 │
│ │ │ │ LoopDetector ④检测 │
│ │ │ │ SUB_ITEM_SEARCH_CAP=1 → 搜满硬挡 │
│ │ │ └─ 返回截断后候选 JSON │
│ │ │ 子 #2..#5: 各搜各的平台 │
│ │ └──────────────────────────────────────────┘
│ └─ 合并 5 份候选
└─ 后续由主 loop 精挑/比价/收尾plaintext1.3 分步详解#
fork 三件事判断(满足任一即 fork):① 能并行 ② 要隔离 ③ 链够深。追问/闲聊/单次工具调用/收尾不 fork。
同质的”相同”与”不同”
相同:FULL_TOOL_SET(13 个,同一 list 引用)、system prompt、中间件栈——能力同质。
不同:① LLM 档位——子用 get_fast_llm() 关 reasoning,主用 get_llm()。② thread_id——sub-{uuid[:8]}-d{depth} 隔离对话历史。③ session_dir——继承父的。④ 执行层权限——depth_gate 在 depth≥1 硬拦聚合/终结工具(price_compare、shopping_summary 等)和上下文工具(planner、category_insight)。能力同质但授权不同质。
延迟求值解循环依赖
make_dispatch_tools(lambda: FULL_TOOL_SET) 传 lambda 而非直接传列表——否则传进去的是 extend 前的快照,子 Agent 的工具集缺 dispatch_tool 自身。
三层并发控制
| 层 | 控制什么 | 机制 | 原语 |
|---|---|---|---|
| ① fork 预算 | 总共 fork 几轮 | ForkBudget(max_parallel=1) | ContextVar 存可变对象引用 |
| ② 并发度 | 同时几个子在跑 | Semaphore(5) | ContextVar 同引用 |
| ③ 树级检索 | 全树累计检索次数 | TREE_RETRIEVAL_BUDGET=8 | module dict(key=session_dir) |
③ 用 dict 不用 ContextVar:asyncio 创建子任务拷贝 context,子 set 父读不到。
异常隔离
_run_sub_agent() 所有异常 catch 转字符串:ForkLimitExceeded → 拒绝、TimeoutError → 超时、GraphRecursionError → 迭代超限。asyncio.gather(return_exceptions=True) 一个子炸了不连累其余。
1.4 技术选型#
| 组件 | 选了什么 | 为什么 | 何时换 |
|---|---|---|---|
| 隔离原语 | ContextVar | asyncio Task 级独立副本 | 多进程 → multiprocessing 级 |
| 聚合原语 | module dict + session_dir | 最轻量跨 Task 共享 | 多机 → Redis |
| fork 模型 | 同质(全工具集 + 权限闸) | 维护一份不维护 N 份;保 cache 前缀 | 子任务差异极大 → 专用化 |
| 平台覆盖 | _ensure_platform_coverage 自动补齐 | 弱模型不听”覆盖全部 5 平台” | 只用强模型 → 可去掉 |
二、踩坑实录#
坑 1:ContextVar 做全树计数器,子写了父读不到#
- 5 个子各搜 3 次,父计数器 0。改 module dict 按 session_dir 聚合。隔离用 ContextVar,聚合用 module dict。
坑 2:同质不只是工具——指令也必须一致#
- 子拿着”鼓励 fork”的正式 prompt 但手里只有玩具工具。改成
make_dispatch_tools接受 system_prompt 参数。
坑 3:recursion_limit 是 super-step 数#
- 设 6 只跑了 3 轮。改成
MAX_ITERATIONS * 2 + 1。
坑 4:弱模型 prompt “覆盖全部 5 平台”约等于白说#
- 关 reasoning 后只列 3 条 demands。加
_ensure_platform_coverage()机制补齐。
三、验收与量化#
| 场景 | 验证点 |
|---|---|
| 递归 fork 被拦 | 子调 dispatch_tool → ForkLimitExceeded |
| 子超时 | 90s → TimeoutError 转字符串 |
| 结果截断 | >4000 token 尾部截断 |
| 循环检测 | 窗口内同工具 4 次 → 换策略提示 |
| 异常不崩主 loop | 子异常转字符串,主正常继续 |
| 平台自动补齐 | 只写 3 条 → 补到 5 条 |
| 跨 fork 检索计数 | 5 个子各搜 1 次 → 全树计数=5 |
四、面试问答#
Q1: 为什么不 asyncio.gather(5 × item_search) 就完了?#
每个平台检索不只一次 item_search 调用。子可能换关键词重搜、查品类常识校准——它自己也是完整 Agent 循环。
Q2: 子和父到底什么不同?#
四点:LLM 档位(关 reasoning)、thread_id 独立、session_dir 继承、执行层权限(depth_gate 拦聚合/终结工具)。能力同质但授权不同质。
Q3: 为什么讨论过”确定性并发”但保留 fork?#
确定性 asyncio.gather(5 × item_search) 确实更省 token(不需要子 Agent 的模型调用开销)。但保留 fork 因为:① 教学架构展示可组合性 ② 子遇到意外能灵活应对 ③ fork 的安全四层是项目核心展示点。
Q4: 检索预算为什么不用 ContextVar?#
asyncio.create_task 拷贝 context,子 set 父读不到。改 module dict + session_dir 做 key。隔离用 ContextVar,共享用 module dict。
Q5: ForkBudget 用 ContextVar 存可变对象引用为什么能跨 Task 共享?#
ContextVar 拷贝的是引用不是值。父和子拿的是同一个 ForkBudget 对象的引用,charge/check 操作修改的是同一个对象。只要只有主 loop charge(子不 charge),就不存在竞态。
五、前沿概念#
5.1 同质 fork vs 专用化 Multi-Agent#
同质:子是主的完整克隆,一份维护全局生效。专用化:每类子 Agent 有定制工具集和 prompt,更精准但维护 N 份。选型看子任务差异度——跨平台检索差异小选同质,NLP+图像差异大选专用。
5.2 asyncio 并发原语选型#
ContextVar:Task 级独立副本,适合隔离场景(thread_id、fork_depth)。module dict:进程内全局共享,适合聚合场景(检索预算、候选登记表)。Semaphore:控制并发度,排队不拒绝。Event/Lock:同步点,本项目用于 fork 内部协调。
六、诚实边界#
| 维度 | 做了 | 没做 |
|---|---|---|
| fork 深度 | MAX_FORK_DEPTH=1(主→子) | 主→子→孙(两层以上) |
| 并发控制 | 进程内 Semaphore + ForkBudget | 多进程/多机需 Redis |
| 检索聚合 | module dict 进程内跨 Task | 多 worker 各算各的 |
| 子 Agent 模型 | 统一用 get_fast_llm() | 无 per-task 模型选择 |
| 教学 vs 效率 | 保留 fork(教学价值) | 确定性并发更省 token |