一、引言:单 Agent 的四道天花板
过去一年,大多数人对 LLM 的用法还停留在“一个对话框解决所有问题”。
但只要任务稍微复杂一点——写一份竞品调研报告、排查一个跨模块的线上故障、把一个需求拆成可合并的代码——单 Agent 架构就会撞上四道天花板:
第一,上下文窗口的物理限制。一个 Agent 既要记住原始目标,又要容纳十几个网页的原文、几十轮工具调用结果,还要保留自己的推理链。窗口再大也会被撑爆,而一旦触发截断,早期目标就被“遗忘”,模型开始跑偏。更糟的是,这种跑偏是静默的:你不会收到任何报错,只会看到一份看起来完整、实际上和原始需求已经偏离的报告。很多“模型变笨了”的抱怨,本质上是上下文管理失败,而不是模型能力退化。
第二,角色冲突。“你既是严谨的调研员,又是富有创造力的文案,同时还要做事实核查”——这种混合系统提示词会让模型在相互矛盾的偏好之间摇摆,最终每一项都做得平庸。调研员需要克制、只写有出处的结论;文案需要发散、追求表达张力;核查员需要怀疑一切、专门挑刺。当这三套价值观被塞进同一个 system prompt,模型会在生成每一句话时同时受到三种力量的拉扯,结果就是“既不严谨也不出彩”。这不是提示词写得不够好,而是任务本身就不该由一个角色承担。
第三,串行瓶颈与错误累积。单 Agent 只能一步步来,第 7 步的幻觉会被第 8 步当作事实继续推理,错误呈指数级传播。想象一条 12 步的流水线,每步 95% 的准确率看似很高,但 0.95¹² ≈ 54%,整条链路已经接近抛硬币。而多 Agent 架构允许在关键节点插入独立的校验角色,把错误在传播路径上就地掐断。
第四,无法并行。十个子问题本质上互不依赖,却没理由排队执行。假设每个子问题要 20 秒,串行就是 200 秒,而并行只需 20 秒出头——这不是优化,是数量级的差距。更关键的是,串行执行时后一个任务会“看到”前一个任务的完整推理痕迹,容易被带偏;并行执行天然实现了上下文隔离。
多 Agent 协作要解决的核心问题,就是上下文隔离 + 任务分解 + 结果归并。2024 年 Anthropic 在 Multi-Agent Research System 的工程博客里给出的数据很能说明问题:采用 Orchestrator-Worker 结构后,在广度优先的调研任务上,相比单 Agent 有约 90% 的性能提升。当然,代价也很明确——同一篇博客提到,多 Agent 系统的 token 消耗约为普通对话的 15 倍。所以这不是“无脑上多 Agent”,而是“在任务足够复杂、且可并行时才值得上”。本文不依赖 AutoGen、CrewAI 这类重型框架,用纯 Python +asyncio从零搭一套可运行、可调试的多 Agent 系统。
二、核心原理:多 Agent 协作的四根支柱
2.1 角色即上下文边界
Agent 的本质不是“人格”,而是一组隔离的上下文 + 固定的系统提示词 + 受限的工具集。给每个 Agent 分配角色,真正起作用的是它只看见与自己相关的那部分信息。上下文隔离本身就是最有效的提示工程。
这里有一个反直觉的结论:如果你给两个“角色不同但上下文相同”的 Agent,它们的行为差异会远小于你的预期。因为模型能看到的信息完全一致,角色提示词不过是一层薄薄的滤镜。真正让多 Agent 产生质变的,是信息不对称——调研员看不到写作者的草稿,审稿人看不到调研员的原始推理链,只看到结论。信息在传递过程中被强制压缩,噪声被过滤,幻觉的传播路径被切断。
在工程实现上,上下文隔离有两种常见做法。一种是逻辑隔离:所有 Agent 共用一个进程,但每个 Agent 每次调用都构造全新的messages数组,互不追加。这是本文采用的方式,轻量、可控、便于调试。另一种是物理隔离:每个 Agent 是一个独立的会话线程(Thread)、独立的容器甚至独立的服务,适合需要长期记忆或不同工具权限的场景。绝大多数任务用逻辑隔离就够了。
2.2 任务分解:从目标到 DAG
Planner 接收用户目标,输出一组带依赖关系的子任务。用有向无环图(DAG)表达依赖,而不是简单的列表,是因为真实任务里“汇总”必须等待“调研”完成,“对比”必须等待“两个方案各自成型”。DAG 让调度器可以最大化并行度。
如果只用列表,你实际上默认了“所有任务按顺序执行”,这会把并行度压到 1;如果只用集合,你又丢掉了依赖信息,调度器不知道谁必须先跑。DAG 是这两者的最小完备表达:deps字段既告诉调度器“等谁”,也告诉下游 Agent“该带哪些上游产物进上下文”。
DAG 还有一个隐性好处:它让失败可定位、可重试。某个节点失败了,你可以只重跑该节点及其下游,而不是整个流程从头发起。在第 5 节我们会看到,这一点在生产环境里非常值钱。
2.3 共享状态:黑板模式
Agent 之间不直接对话,而是通过一块**黑板(Blackboard)**交换产物。黑板只存最终产物,不存推理过程。这样做的关键收益是:上游 Agent 冗长的思维链不会污染下游,传递的是压缩后的结论。
“黑板模式”并不是新概念,它来自上世纪 70 年代的经典 AI 架构(Hearsay-II 语音理解系统),核心思想是:多个独立的知识源围绕一块共享数据结构协作,谁有能贡献的部分就写入,谁需要就读取。大模型时代的黑板,只不过是把这套机制重新用了一遍。
工程上,黑板必须承担三件事:
- 压缩:写入前对内容做长度控制,超长就截断并标注,防止下游上下文爆炸。
- 寻址:下游 Agent 通过
deps精确读取自己需要的条目,而不是把整块黑板全塞进去。 - 隔离:写入的是结论,不是过程。如果一个 Worker 输出了三页推理链,中间还有反复推翻的段落,直接传给下游只会让下游困惑。
2.4 调度与汇总
调度器负责拓扑排序、并发控制、失败降级;汇总器(Aggregator)负责去重、消解矛盾、按目标重新组织。这本质上是一次 Map-Reduce:Map 阶段并行生产片段,Reduce 阶段串行收敛。
调度器的三个关键词值得展开:
- 拓扑排序:按依赖关系决定谁先谁后,常见算法是 Kahn(基于入度)。调度器可以每完成一批就重新计算“当前就绪”的节点,从而把并行度吃满。
- 并发控制:不是所有任务都该同时跑。并发过高会撞上 API 速率限制,也会让 token 账单失控,所以需要
Semaphore限流。 - 失败降级:一个 Worker 挂了,是整单失败,还是用占位内容继续?这取决于任务重要性。通常建议“关键路径失败即终止,非关键路径失败则降级并在最终报告中标注缺口”。
┌──────────────┐ 用户目标 ───▶ │ Planner │ ──▶ DAG(Task[]) └──────────────┘ │ ┌───────────────┼───────────────┐ ▼ ▼ ▼ ┌─────────┐ ┌─────────┐ ┌─────────┐ │ Worker1 │ │ Worker2 │ │ Worker3 │ ← 并行执行 └────┬────┘ └────┬────┘ └────┬────┘ └───────────────┼───────────────┘ ▼ ┌──────────────┐ │ Blackboard │ ← 只存产物,不存思维链 └──────┬───────┘ ▼ ┌──────────────┐ │ Aggregator │ ──▶ 最终报告 └──────────────┘三、实战:从零实现“技术调研报告生成器”
3.1 最小 Agent 内核
先写最底层:一个无状态、可复用的 Agent 类。它只做一件事——把系统提示词、上游资料、当前指令拼成消息列表,调一次模型。
# agents/core.pyimportosfromdataclassesimportdataclass,fieldfromopenaiimportAsyncOpenAI client=AsyncOpenAI(api_key=os.environ["OPENAI_API_KEY"],base_url=os.getenv("OPENAI_BASE_URL","https://api.openai.com/v1"),timeout=60.0,max_retries=3,# SDK 内置指数退避)@dataclassclassTask:id:strgoal:strdeps:list[str]=field(default_factory=list)result:str|None=NoneclassAgent:"""最小可用 Agent:一个系统提示词 + 一次无状态调用。"""def__init__(self,name:str,system_prompt:str,*,model:str="gpt-4o-mini",temperature:float=0.3,json_mode:bool=False):self.name=name self.system_prompt=system_prompt self.model=model self.temperature=temperature self.json_mode=json_modedef_build_messages(self,instruction:str,context:str)->list[dict]:"""把「角色设定 + 上游产物 + 当前指令」拼成消息列表。 关键约定:上游产物永远放在 user 消息里,绝不放进 system。 system 是角色的“宪法”,必须稳定、不被数据污染—— 否则上游内容里的一句“忽略之前的指令”就可能劫持当前 Agent。 """parts=[]ifcontext.strip():parts.append(f"<参考资料>\n{context}\n</参考资料>")parts.append(f"<任务>\n{instruction}\n</任务>")return[{"role":"system","content":self.system_prompt},{"role":"user","content":"\n\n".join(parts)},]asyncdefrun(self,instruction:str,context:str="")->str:"""执行一次调用。无状态:每次都是全新的消息数组。"""resp=awaitclient.chat.completions.create(model=self.model,temperature=self.temperature,messages=self._build_messages(instruction,context),response_format={"type":"json_object"}ifself.json_modeelseNone,)returnresp.choices[0].message.contentor""这段代码里有三个值得说明的设计选择。
第一,Agent 是无状态的。它不保存对话历史,每次run都从零构造消息数组。这是刻意的:有状态 Agent 会随调用次数增长而上下文膨胀,也会让并发变得危险(同一个 Agent 被两个任务同时调用,历史会串味)。把“状态”统一交给外部的黑板管理,Agent 只负责“给定输入,产出输出”,职责单一,单元测试也好写。
第二,json_mode是结构化输出的开关。Planner 必须返回严格 JSON,所以它开启这个模式;普通的写作型 Agent 则关闭。开启后,模型会被约束返回合法 JSON 对象,大幅降低解析失败率——注意它保证的是“合法 JSON”,不保证“符合你的 schema”,所以后面还要自己做字段校验。
第三,超时与重试交给 SDK。timeout=60.0防止单个请求无限挂起,max_retries=3让 SDK 在网络抖动或 429 限流时自动指数退避重试。不要自己写while True: try: ...,那通常会把 429 变成雪崩。
3.2 Planner:把目标拆成 DAG
Planner 是整个系统的入口,也是质量的上限。它拆错了,后面全错。所以给它的提示词必须极其具体,明确约束“可执行、可并行、无环、必须收口”。
# agents/planner.pyimportjsonfrom.coreimportAgent,Task PLANNER_PROMPT="""你是一个任务规划器,唯一职责是把用户目标拆成 3-6 个可独立执行的子任务。 硬性规则: 1. 每个子任务的 goal 必须自包含,只看 goal 就能执行,禁止出现“同上”“参考上文”这类指代。 2. 有先后依赖时用 deps 声明,值为上游任务的 id。 3. 依赖图必须无环。 4. 必须包含且只包含一个“汇总”任务,其 deps 覆盖其余所有任务。 5. 只输出 JSON,不要任何解释、不要 markdown 代码块。 输出格式: {"tasks": [{"id": "t1", "goal": "...", "deps": []}, ...]} """classPlanner(Agent):def__init__(self,**kwargs):super().__init__(name="planner",system_prompt=PLANNER_PROMPT,temperature=0.2,# 规划要稳定,压低随机性json_mode=True,**kwargs,)asyncdefplan(self,goal:str)->list[Task]:raw=awaitself.run(f"用户目标:{goal}")data=json.loads(raw)# json_mode 已保证是合法 JSONtasks=[Task(id=item["id"],goal=item["goal"],deps=item.get("deps",[]))foritemindata["tasks"]]self._validate(tasks)returntasks@staticmethoddef_validate(tasks:list[Task])->None:"""模型会撒谎,依赖图必须自己校验。"""ids={t.idfortintasks}iflen(ids)!=len(tasks):raiseValueError("存在重复的任务 id")fortintasks:fordint.deps:ifdnotinids:raiseValueError(f"任务{t.id}依赖了不存在的{d}")# Kahn 算法检测环indeg={t.id:len(t.deps)fortintasks}children={t.id:[]fortintasks}fortintasks:fordint.deps:children[d].append(t.id)queue=[tidfortid,ninindeg.items()ifn==0]seen=0whilequeue:cur=queue.pop()seen+=1fornxtinchildren[cur]:indeg[nxt]-=1ifindeg[nxt]==0:queue.append(nxt)ifseen!=len(tasks):raiseValueError("依赖图中存在环,无法完成拓扑排序")_validate不是可选项。实测中,即便用了 JSON 模式,模型仍有一定概率产出指向不存在 id 的deps,或者在多轮迭代后生成环。让它在进入调度器之前就爆掉,比让调度器死锁好得多——前者的报错信息清晰,后者你只会看到进程卡住。
3.3 Blackboard:只存结论,不存过程
# agents/blackboard.pyclassBlackboard:"""Agent 之间唯一的通信介质。只存最终产物,不存思维链。"""def__init__(self,max_chars_per_item:int=4000):self._store:dict[str,str]={}self._max=max_chars_per_itemdefput(self,task_id:str,content:str)->None:content=(contentor"").strip()ifnotcontent:content="(该子任务未产出有效内容)"iflen(content)>self._max:content=content[:self._max]+"\n...[内容超长已截断]"self._store[task_id]=contentdefget(self,task_id:str)->str:returnself._store.get(task_id,"")defcontext_for(self,deps:list[str])->str:"""只把当前任务声明的依赖项拼进上下文,而非整块黑板。"""blocks=[]fordindeps:val=self.get(d)ifval:blocks.append(f"### 来自任务{d}的结论\n{val}")return"\n\n".join(blocks)defall_results(self)->dict[str,str]:returndict(self._store)context_for是黑板的灵魂。它实现了“按需投喂”:一个任务只看到它deps里列出的上游产物。如果不做这一步,而是把黑板全量塞给每个 Agent,那么并行带来的上下文隔离优势瞬间归零,你又回到了单 Agent 撑爆窗口的老路。max_chars_per_item则是第二道保险——单个产物的长度上限,防止某个 Agent 输出了两万字长文把下游拖垮。
3.4 Scheduler:拓扑排序 + 并发控制
# agents/scheduler.pyimportasynciofrom.coreimportTaskclassScheduler:def__init__(self,blackboard,worker_for,*,max_concurrency:int=4):self.board=blackboard self.worker_for=worker_for# Callable[[Task], Agent]self.sem=asyncio.Semaphore(max_concurrency)asyncdef_execute(self,task:Task)->None:asyncwithself.sem:# 全局限流,保护 API 配额ctx=self.board.context_for(task.deps)agent=self.worker_for(task)output=awaitagent.run(task.goal,ctx)self.board.put(task.id,output)asyncdefrun(self,tasks:list[Task])->None:pending={t.id:tfortintasks}done:set[str]=set()running:dict[asyncio.Task,str]={}whilependingorrunning:# 1) 把所有依赖已满足的任务挂上事件循环fortidinlist(pending):task=pending[tid]ifall(dindonefordintask.deps):fut=asyncio.create_task(self._execute(task))running[fut]=tiddelpending[tid]ifnotrunning:raiseRuntimeError(f"依赖图无法推进,剩余:{list(pending)}")# 2) 等最快的一个完成,立刻回到循环重新找就绪节点finished,_=awaitasyncio.wait(running,return_when=asyncio.FIRST_COMPLETED)forfutinfinished:tid=running.pop(fut)fut.result()# 有异常在此抛出,快速失败done.add(tid)调度循环的心跳是 `asyncio.
更多硬核网安与AI工具包,请扫码获取完整源码!
wait(…, FIRST_COMPLETED):每有一个任务完成,就重新扫描一遍pending`,把