多个 Agent 并发跑完结果对不上,到底该不该等齐所有分支
数据截至 2026-07,各产品的额度与报错口径以官方最新说明为准。
并发编排出问题的时候,绝大多数人第一反应是”模型不稳定”,然后去换模型、加重试、调温度——但真正的病因通常是等待语义写错了:你以为在等所有分支,实际上代码只等了先返回的那几个;或者反过来,明明可以先走的环节被一个不相干的分支拖成了串行。 这类问题的特征很好认:单独跑每个分支都对,合起来跑就时对时错,重跑一次结果还不一样。模型是无辜的,是你的汇聚点没定义清楚。
这篇只谈编排结构本身。站内另外两篇分工不同:Agent 角色分工讲的是把一个大任务拆成几个角色、每个角色负责什么,是”拆”的问题;Agent 日常运维讲的是跑起来之后的监控、告警和值班,是”守”的问题。本篇夹在中间,讲的是拆完之后这几路怎么并、在哪个点合、合不上怎么办——是”接”的问题。
一、先分清你手上是哪种并发结构
排查之前先给结构归类,因为三种结构的失败长相完全不同,混着排会绕远路。
扇出扇入(fan-out / fan-in):一个输入分裂成 N 个同构分支,各自独立产出,最后汇总成一个结果。典型场景是让多个 Agent 分别读代码库的不同目录再合并成一份评估,或者对同一个问题跑多种方案再择优。它的特点是分支之间没有依赖,汇聚点是唯一的同步点。
屏障(barrier):多个异构任务必须全部到齐,下一阶段才能开始。比如生成接口定义、生成数据模型、生成迁移脚本三件事同时做,但只有三份都在手上才能做一致性校验。屏障和扇入的区别在于:扇入的产物是”合并”,屏障的产物是”许可”——屏障后面那步需要的是完整上下文,缺一个就没法判断。
流水线(pipeline):任务分成有序的几段,每段处理完就把结果推给下一段,多个任务在不同段上同时流动。比如批量处理 200 个文件,读取、改写、校验三段各自并发,中间用队列衔接。流水线里根本不存在”等齐”,等齐就是把流水线退化成了分批串行。
判断口诀:下游需要全体信息才能做决策的,必须等齐;下游只是把结果堆起来的,可以边到边收;下游只关心单条记录的,压根不该等。 大量并发 bug 来自把第三种当第一种写。
二、按现象定位成因
下面这张表是我实际排这类问题时的主路径。先看现象,再挑验证手段,别上来就改代码。
| 现象 | 大概率成因 | 怎么验证 | 处置动作 |
|---|---|---|---|
| 每次重跑结果不同,但单跑每个分支都正确 | 汇聚点用了”先到先得”语义,实际只收了部分分支 | 在汇聚前打印已收到的分支 ID 集合与期望集合做差集 | 改成显式等待全集合,或明确声明这是竞速语义并记录被丢弃的分支 |
| 结果顺序乱、拼接后的文档段落错位 | 按完成时间收集而非按分支序号写回 | 给每个分支带序号,汇聚后检查序号是否连续递增 | 用索引写回预分配的结果数组,禁止 append |
| 整体耗时约等于所有分支耗时之和,而不是最慢那一个 | 名义并发实际串行,通常是共享了一个带锁的客户端或在循环里 await | 记录每个分支的开始时间戳,看是否首尾相接 | 把 await 移出循环,先收集任务再统一等待 |
| 少数分支永远不返回,整体卡死 | 无超时的等待 + 上游连接静默断开 | 检查是否有 ETIMEDOUT/ECONNRESET 被吞掉;给等待加一个明确上限 | 给每个分支设独立超时,超时按缺失处理而非无限等 |
| 大量分支同时返回 429 | 扇出宽度超过了服务端并发或速率限制 | 观察 429 是否集中在扇出瞬间、随并发数线性增加 | 加信号量控制在途数量 + 指数退避重试 |
| 汇聚结果里出现别的分支的内容 | 分支之间共享了可变上下文对象 | 给每个分支塞入唯一标记词,看汇聚结果里有没有串味 | 每分支深拷贝独立上下文,禁止共享可变结构 |
| 部分分支返回空但流程照常”成功” | 失败被捕获成空值,汇聚点没做完整性校验 | 统计每个分支产物的长度/条数分布,看有没有 0 | 汇聚前加断言,缺失就显式失败或降级,不许静默通过 |
| 中途改动的文件被另一个分支覆盖 | 多个分支写同一路径,没有隔离工作区 | 看这些分支的写入路径是否有交集 | 分支级隔离(各自独立目录或工作树),最后再合 |
表里最后一条在代码类任务里特别常见。多个 Agent 并行改同一个仓库,谁后写谁赢,前面那份工作就凭空消失了。冲突本身怎么收拾,多会话并发冲突那篇讲得更细,本篇不重复展开。
三、什么时候必须等齐,什么时候不必
这是本篇的核心判断,我给三条可执行的准则。
准则一:下游要做全局决策,就必须等齐。 判断依据是问自己”如果少一个分支的结果,下游给出的答案会不会错,而不只是不完整”。会错的,就是硬屏障。比如你要从五个候选方案里选最优,缺一个候选,选出来的”最优”就是错的结论而非部分结论。一致性校验、去重、排序、投票、择优,全都属于这一类。
准则二:下游只做累加,就边到边收。 汇总报告、日志归集、把每个模块的说明拼成一份文档,这些场景下少一段就是少一段,不会让已有内容变错。这时候用流式收集,先到先写,同时记录缺失清单,最后单独补。这样长尾分支不会拖住整体,用户也能提前看到部分结果。
准则三:能容忍缺失的,提前定义降级结果。 屏障最怕”必须等齐”和”有分支永远不来”同时成立。破解方法是把硬屏障改成”带兜底的软屏障”:设一个明确的等待上限,到点未到的分支用一个预先定义好的降级值代替,并在最终产物里标注哪些位置是降级的。关键是降级值必须是显式的、可识别的,而不是空字符串——空值会一路往下渗透,最后你根本找不到问题起点。
还有一个容易被忽略的点:等齐的粒度可以下调。 很多人默认屏障必须放在全部 N 个分支之后,其实可以分组。比如 20 个分支按主题分成 4 组,组内等齐、组间流水,既保证了组内的一致性判断,又不至于被某一组的慢分支拖垮全场。这个改法通常能把总耗时砍掉一大截,而语义上没有任何损失。
四、动手:怎么改
按上面定完性质,改法就很确定了。这里给几段通用代码,都是语言标准库能力,不涉及任何产品专有接口。
严格等齐(硬屏障)在 Python 里是这样:
import asyncio
async def fan_in_strict(tasks):
# return_exceptions=True 保证一个分支炸掉不会让其余分支被取消
results = await asyncio.gather(*tasks, return_exceptions=True)
failed = [i for i, r in enumerate(results) if isinstance(r, Exception)]
if failed:
raise RuntimeError(f"分支未全部完成,缺失索引: {failed}")
return results
注意 return_exceptions=True 这个细节。默认值是 False,第一个异常抛出后其余任务会被取消,你不仅丢了失败的那个,还丢了本来能成功的那些,日志里只剩一条异常,根本看不出到底是谁的问题。先全收,再判断,最后统一报错,排查成本差一个数量级。
带上限的软屏障:
async def fan_in_soft(coros, timeout_s, fallback):
# asyncio.wait 不接受裸协程(较新版本会直接抛 TypeError),必须先包成 Task
tasks = [asyncio.ensure_future(c) for c in coros]
done, pending = await asyncio.wait(tasks, timeout=timeout_s)
for t in pending:
t.cancel()
out = []
for t in tasks:
if t in done and not t.cancelled() and t.exception() is None:
out.append(t.result())
else:
out.append(fallback) # 显式降级标记,下游可识别
return out
这段有两个容易写错的地方。第一,asyncio.wait 和 asyncio.gather 的入参要求不一样:gather 可以直接吃协程对象,wait 要求传进来的是 Task 或 Future,把协程直接丢给它会报参数类型错误。所以这里先统一 ensure_future 包一层,顺便让后面能按同一批对象做顺序还原。第二,判断成功要先排除被取消的任务——被取消的 Task 同样算”已完成”,会躺在 done 集合里(外部提前取消、或取消与超时同刻竞态都会出现),而对已取消的 Task 调 exception() 抛的是 CancelledError 而不是返回 None。写成 not t.exception() 就会在本该安静降级的那条路径上,炸出一个跟业务毫无关系的异常。
另外这里刻意保持了输出与输入的顺序一一对应,而不是用 done 集合直接拼——done 是无序的,用它拼出来的结果顺序不可预测,这正是第二节表里”段落错位”那一行的成因。
控制扇出宽度,避免自己把自己打成 429:
async def run_all(jobs, width=4):
sem = asyncio.Semaphore(width) # 在途上限按服务端限制调,各家规则不同且会调整
async def guarded(coro_fn, *args):
async with sem:
return await coro_fn(*args)
return await asyncio.gather(*(guarded(fn, *a) for fn, a in jobs))
信号量在协程里创建、不放模块顶层,是为了避免它跟事件循环的绑定时机在不同运行环境下产生差异;在异步入口里现建现用永远是安全的写法。
至于具体开多宽,没有通用答案。各家服务的并发与速率规则不同且会随时调整,以官方最新说明为准。可行的做法是从小往大试,观察 429 出现的临界点,然后取临界值的七八成作为常驻值。
并发场景一定要给每个分支独立的重试预算,而不是在汇聚层统一重试。 汇聚层重试意味着一个分支失败就要把全部 N 个分支重跑一遍,成本乘以 N,而且已经成功的分支再跑一次结果可能还不一样,反而引入新的不一致。重试要贴着最小失败单元做,这一块的细节在 Agent 失败重试里讲得更细。
五、什么情况下别再折腾了
并发编排有一个很难受的特点:它的收益上限是固定的(最多把总时长压到最慢分支那么长),但复杂度的上升是没有底的。所以要提前定好止损点。
止损信号一:并发版本的正确率低于串行版本。 只要出现这个情况,立刻退回串行,别调了。并发是性能优化,性能优化不能拿正确性换。判断方法是拿同一批输入串行跑三遍、并发跑三遍,两组分别看内部一致性。这里必须先扣掉模型本身的采样随机性——串行三遍就已经不一致的部分,那是生成侧的波动,不归编排管;只有”串行三遍稳定、并发三遍飘”这种差值,才指向编排层有没定位到的共享状态或竞态。想把这个差值看干净,先把温度类采样参数压到最低、固定住能固定的随机源,再做对比。确认是编排层的问题之后,继续加锁加重试只会把问题埋得更深,退回串行才是止损。
止损信号二:为了并发而引入的协调代码超过了业务代码的体量。 当你发现自己在写分支状态机、写补偿逻辑、写部分失败的回滚,说明这个任务的耦合度根本不适合并发拆分。这时候正确的动作是回到拆分层重新划边界,而不是在编排层硬扛。
止损信号三:加宽并发后总耗时不降反升。 这通常意味着瓶颈不在你这一侧——可能是服务端排队,可能是本地 IO 或内存,也可能是重试风暴。继续加宽只会让 429 更密集。验证办法是记录”提交到首字节”的时间,如果这个值随并发数上升,瓶颈就在下游,加宽没有意义。
回滚点怎么留: 编排代码要保证串行路径永远可用,用一个开关切换,不要把串行分支删掉。我见过太多项目并发改造做到一半上线出事,想退回去发现串行代码早就被删了,只能硬着头皮修。留着串行路径的成本极低,收益是你随时有一条确定能用的路。
换条路的判断依据: 如果任务本身是长尾分布——大部分分支很快、极少数分支特别慢——那么屏障式的等齐永远不划算,应该换成流水线加异步回填:先出一版不含慢分支的结果交付,慢分支完成后再补进去。差别在于首次可见时间的口径变了:屏障式等齐,用户看到任何东西都要等到最慢那个分支返回;改成回填之后,看到八成内容的时间取决于第八十百分位那个分支,剩下两成再异步补。分布越长尾,这两个数差得越远——这也是为什么该不该改成回填,要拿分支耗时的分位数说话,而不是拿平均值。平均值会被少数极慢分支拉高,反而看不出长尾有多长。
六、避坑清单
坑一:把 Promise.race / asyncio.wait(..., return_when=asyncio.FIRST_COMPLETED) 当成 Promise.all / gather 用。 为什么会踩:这两组 API 名字接近、签名接近,改代码的时候顺手替换很难在 review 里看出来,而且在分支都很快的测试环境里表现一样。怎么避:在汇聚点强制断言收到的结果条数等于扇出条数,测试用例里必须有一个人为拖慢的分支。
坑二:在循环里 await,写出假并发。 为什么会踩:for x in items: await f(x) 读起来非常自然,形式上还有 await 看着像异步。怎么避:立一条规矩——循环体里不许出现 await,只许收集任务对象,出了循环再统一 gather。这条规矩机械但有效,review 时一眼能扫出来。
坑三:多个分支共享同一个上下文/消息列表对象。 为什么会踩:为了省内存或者图省事,把同一个上下文传给所有分支,而 Agent 通常会往上下文里追加内容,于是分支之间互相污染。怎么避:分发前深拷贝,或者干脆用不可变结构。验证手段是给每个分支塞一个唯一标记词,跑完 grep 结果里有没有别人的标记。
坑四:并发分支往同一个文件写。 为什么会踩:每个分支单独看逻辑都对,“写结果到 output.md”这种代码在串行时代写了无数遍,改并发时没人想起它。怎么避:分支的产出一律先留在内存或独立临时文件里,由汇聚层唯一负责落盘。代码改动型任务则用独立工作区隔离。
坑五:失败被 try/except 吞成空值。 为什么会踩:为了让”一个分支挂了不影响整体”,顺手写了个宽泛的异常捕获返回空。结果流程照常走完,产物少了一块,而没有任何报错。怎么避:捕获可以,但必须返回一个可识别的失败对象(带分支 ID 和原因),并且汇聚层要统计失败数并输出到日志摘要里。静默成功比明确失败危险得多。
坑六:重试放在错误的层级。 为什么会踩:在最外层加重试是最省事的写法,一行装饰器搞定。怎么避:明确一条原则——重试永远贴着最小可独立失败的单元。外层只做整体失败的上报,不做重试。
坑七:并发数写死在代码里,跨环境直接炸。 为什么会踩:本地调试时开 16 路很爽,上到共享环境或者换了服务端配额就开始密集 429。怎么避:并发数一律从环境变量读,给一个保守的默认值。下面这个变量名是自己起的、不是任何产品认的约定,照抄名字没意义,照抄做法才有意义。
export AGENT_FANOUT=4
坑八:把成本当成常量。 为什么会踩:扇出 N 路意味着 token 消耗大致乘以 N,而串行时形成的成本直觉完全失效。并发跑一次的花费可能远超预期。怎么避:并发编排上线前先加计量,按分支维度统计消耗。具体做法见 Agent 成本失控。
坑九:用海外工具做并发压测时忽略区域限制。 部分海外 AI 工具与模型服务,官方明确对中国大陆有区域限制、不支持直连,你在本地测出来的并发表现不能代表实际可用状态。市面上确实存在第三方中转,但稳定性和合规性都需要自己评估,我不做推荐也不给渠道。做技术选型时,先确认服务在你的部署地域是否可正常访问,再谈并发调优。
收束
并发编排的排查顺序其实很固定:先归类结构(扇出扇入 / 屏障 / 流水线),再按现象查表定因,然后按”下游是否需要全局信息”决定等不等齐,最后给每个等待点配上限和降级值。 顺序反过来做——比如上来就调并发数、加重试——基本都是在给自己制造新的不确定性。
上线前过一遍这份自检清单:
- 汇聚点是否显式断言了”收到条数 == 扇出条数”?
- 结果是按索引写回的,还是按完成顺序 append 的?
- 每个分支是否有独立超时?超时后的降级值是不是可识别的?
- 分支之间有没有共享可变对象、有没有写同一个路径?
- 重试贴在最小失败单元上了吗?还是套在最外层?
- 并发数是从配置读的吗?串行回滚路径还留着吗?
- 有没有按分支统计消耗和失败数,落到日志摘要里?
这七条里但凡有一条答不上来,就先别急着加宽并发。把等待语义写对,比把并发数调大值钱得多。