1. 多Agent协作到底在解决什么问题
单Agent跑任务,跑到一定复杂度就会撞墙。我最早做自动化流程的时候,一个Agent包揽需求解析、资料检索、代码生成、结果校验,提示词写到三千字,工具挂了十几个,结果就是:它开始"精神分裂"。前面说要按A方案走,中间检索回来一堆信息,它转头就按B方案生成了,最后校验环节又用C方案的标准去检查。这不是模型不行,是架构本身就不对。
多Agent协作要解决的核心问题就三个:上下文污染、职责耦合、单点瓶颈。一个Agent的上下文窗口是有限的,你往里塞的东西越多,它的注意力就越分散,关键指令被淹没的概率就越大。这跟人一样,你让一个人同时干产品、开发、测试、运维,他不是干不了,是干不好,而且一旦某个环节出问题,你根本不知道是哪个环节的锅。
所以多Agent的本质是把一个大而全的模糊任务,拆成多个小而精的明确任务,每个任务交给一个专职Agent,再通过一套调度机制把它们串起来。听起来简单,但真正落地的时候,协作架构怎么设计、任务怎么调度、Agent之间怎么通信、冲突怎么解决,每一个都是坑。
这篇文章我会从协作架构的分类讲起,然后深入到任务调度的具体实现,再给出一套完整的复杂AI协同任务构建方案。适合已经跑通过单Agent、想往多Agent方向进阶的开发者,也适合正在做AI应用架构设计的技术负责人。文章里涉及到的代码和配置,都是可以直接拿去改改就用的。
2. 协作架构的四种主流模式与选型逻辑
2.1 从"谁说了算"来区分架构类型
多Agent协作架构,按决策权的分布方式,可以分成四类。这个分类方式比按技术栈分更实用,因为它直接决定了你后面任务调度怎么写。
中心化架构(Orchestrator模式):有一个主Agent充当调度中心,所有子任务由它分配,所有结果由它汇总。这是最容易上手的模式,也是我推荐新手第一个尝试的架构。它的优势是控制流清晰,出问题容易定位。劣势是主Agent容易成为瓶颈,而且主Agent的上下文压力很大,因为它要记住所有子Agent的状态。
去中心化架构(Peer-to-Peer模式):Agent之间平等通信,没有全局调度者。每个Agent根据自己的状态和收到的消息决定下一步动作。这种架构灵活性高,但调试难度直线上升。我试过用这种模式做一个内容审核流水线,三个Agent互相传递任务,结果出现了两个Agent互相等待对方先行动的死锁。后来加了超时机制才解决。
层级架构(Hierarchical模式):中心化的升级版,主Agent下面还有中间层Agent,每个中间层管理一组执行Agent。适合任务层级深、子任务数量多的场景。比如一个大型代码生成任务,顶层Agent拆成"前端""后端""数据库"三个中间层,每个中间层再拆成具体的模块任务。
混合架构(Hybrid模式):实际生产环境里用得最多的其实是混合模式。核心调度用中心化保证可控性,局部协作允许去中心化提高效率。比如主Agent负责拆解和汇总,但检索类Agent之间可以互相直接调用,不用每次都经过主Agent转发。
选型的时候,我一般看三个指标:任务复杂度、Agent数量、容错要求。任务简单、Agent少于5个,中心化就够了。Agent超过10个、任务有明确层级,考虑层级架构。对实时性要求高、允许局部失败,可以试试混合模式。去中心化我一般不建议在生产环境用,除非你有很强的分布式系统调试能力。
2.2 通信机制的选择:消息传递 vs 共享状态
Agent之间怎么交换信息,这是架构设计里第二个关键决策。主流方案有两种:消息传递和共享状态。
消息传递就是Agent之间直接发消息,像微信聊天一样。A发一条"帮我查一下这个数据",B收到后处理完回一条"结果在这"。这种方式的优势是解耦彻底,每个Agent只需要知道"我该给谁发消息"和"收到消息后怎么处理"。劣势是消息格式需要严格定义,而且消息丢失或乱序的时候处理起来很麻烦。
共享状态是所有Agent读写同一个状态存储,比如一个共享的JSON对象或者数据库。A往里面写一个字段,B读这个字段。这种方式的好处是状态一致性容易保证,不需要处理消息乱序问题。坏处是并发写入需要加锁,而且状态结构一旦设计不好,后期扩展很痛苦。
我自己的经验是:任务流程线性、Agent数量少的时候用消息传递;任务流程有分支、Agent需要频繁读取全局信息的时候用共享状态。实际项目里我经常混用:核心流程用消息传递保证解耦,全局配置和中间结果用共享状态方便查询。
注意:不管用哪种通信机制,一定要给消息或状态字段加上版本号和时间戳。我踩过一次坑,两个Agent同时写同一个状态字段,后写的覆盖了先写的,导致整个流程跑偏。加了版本号之后,冲突检测就简单多了。
2.3 任务调度的核心:DAG还是状态机
任务调度这块,最常用的两种模型是DAG(有向无环图)和状态机。
DAG适合任务依赖关系明确的场景。比如"数据清洗→特征提取→模型训练→结果评估",每个节点是一个Agent任务,箭头代表依赖关系。DAG的优势是可视化好、依赖检查简单、可以并行执行无依赖的节点。劣势是它假设任务流程是固定的,一旦运行中需要动态调整流程,DAG就不太够用了。
状态机适合任务流程会根据中间结果动态变化的场景。比如"如果检索到的资料足够,就直接生成;如果不够,就触发补充检索"。状态机把每个Agent任务定义成一个状态,状态之间的转移条件由业务逻辑决定。这种方式灵活,但状态爆炸的问题需要提前考虑。
我一般建议:流程固定的用DAG,流程动态的用状态机,两者可以结合。比如顶层用状态机控制大流程,每个状态内部用DAG控制子任务依赖。这样既有灵活性,又有可控性。
3. 任务调度的具体实现与核心代码
3.1 调度器的基本结构
一个任务调度器,不管用什么语言写,核心就四件事:任务注册、依赖解析、执行调度、结果收集。我用Python写一个最小可用的调度器,你可以直接拿去改。
import asyncio from dataclasses import dataclass, field from typing import Callable, Any from enum import Enum class TaskStatus(Enum): PENDING = "pending" RUNNING = "running" DONE = "done" FAILED = "failed" @dataclass class Task: name: str func: Callable depends_on: list[str] = field(default_factory=list) status: TaskStatus = TaskStatus.PENDING result: Any = None error: str = None class Scheduler: def __init__(self): self.tasks: dict[str, Task] = {} self.results: dict[str, Any] = {} def register(self, task: Task): self.tasks[task.name] = task def _get_ready_tasks(self) -> list[Task]: ready = [] for task in self.tasks.values(): if task.status != TaskStatus.PENDING: continue deps_done = all( self.tasks[d].status == TaskStatus.DONE for d in task.depends_on ) if deps_done: ready.append(task) return ready async def run(self): while True: ready = self._get_ready_tasks() if not ready: break await asyncio.gather(*[self._execute(t) for t in ready]) async def _execute(self, task: Task): task.status = TaskStatus.RUNNING try: deps_results = {d: self.tasks[d].result for d in task.depends_on} task.result = await task.func(deps_results) task.status = TaskStatus.DONE self.results[task.name] = task.result except Exception as e: task.status = TaskStatus.FAILED task.error = str(e)这段代码的核心逻辑是:每次循环找出所有依赖已完成的待执行任务,并行执行它们,直到没有可执行的任务为止。_get_ready_tasks是依赖解析的关键,它检查每个待执行任务的所有依赖是否都已完成。
实际用的时候,你需要在这个基础上加几个东西:超时控制、重试机制、失败传播策略。超时控制是给每个任务设一个最大执行时间,超了就标记失败。重试机制是失败后自动重试N次。失败传播策略是当一个任务失败时,依赖它的任务是直接跳过还是也标记失败。
3.2 Agent任务的封装方式
调度器有了,接下来要把Agent封装成可调度的任务。一个Agent任务的标准结构包括:输入解析、提示词构建、模型调用、输出解析、结果校验。
class AgentTask: def __init__(self, name, role_prompt, model_client, output_schema=None): self.name = name self.role_prompt = role_prompt self.model_client = model_client self.output_schema = output_schema async def __call__(self, deps_results: dict) -> Any: context = self._build_context(deps_results) messages = [ {"role": "system", "content": self.role_prompt}, {"role": "user", "content": context} ] raw_output = await self.model_client.chat(messages) parsed = self._parse_output(raw_output) if self.output_schema: self._validate(parsed) return parsed def _build_context(self, deps_results): parts = [] for dep_name, result in deps_results.items(): parts.append(f"[来自 {dep_name} 的结果]\n{result}") return "\n\n".join(parts)这里有几个实操细节值得说。_build_context里给每个依赖结果加了来源标记,这是为了让当前Agent知道信息是从哪来的,方便它做判断。_parse_output要做健壮性处理,因为模型输出不一定是纯JSON,可能带markdown代码块标记,也可能有额外的解释文字。我一般用正则先把JSON部分抠出来,再解析。
_validate是输出校验,这个非常重要。多Agent系统里,一个Agent的输出格式错了,后面依赖它的Agent全都会崩。校验不通过的时候,我一般会触发一次重试,把校验错误信息拼回提示词里让模型重新生成。
3.3 并行执行与资源控制
多Agent系统跑起来之后,你会发现瓶颈往往不在模型推理,而在并发控制和资源竞争。同时跑10个Agent,每个都在调API,很容易触发速率限制。同时写共享状态,不加锁就会数据错乱。
并行执行的控制,我一般用信号量(Semaphore)来限制同时运行的Agent数量。比如你API的速率限制是每分钟60次,每个Agent平均调用3次模型,那同时最多跑20个Agent。留点余量,设成15比较稳。
class RateLimitedScheduler(Scheduler): def __init__(self, max_concurrent=5): super().__init__() self.semaphore = asyncio.Semaphore(max_concurrent) async def _execute(self, task: Task): async with self.semaphore: await super()._execute(task)共享状态的并发控制,如果用的是内存字典,Python的GIL能保证单次操作的原子性,但"读-改-写"这种复合操作就不行了。我一般用asyncio.Lock来保护关键区段。如果是多进程或者分布式部署,那就得上Redis的分布式锁或者数据库的行锁。
实操心得:并发数不是越大越好。我试过把并发调到50,结果模型API的响应时间从2秒涨到了15秒,整体吞吐反而下降了。后来做了个简单测试,找到响应时间和并发数的拐点,一般设在拐点前20%的位置最稳。
4. 复杂AI协同任务的完整构建流程
4.1 任务拆解:从模糊需求到可执行DAG
拿到一个复杂任务,第一步是拆解。拆解的质量直接决定了后面所有环节的成败。我用的方法叫"三层拆解法":目标层、能力层、执行层。
目标层是明确最终要交付什么。比如"写一份行业分析报告",交付物是一份报告,包含市场概况、竞争格局、趋势判断三个部分。能力层是完成这个交付物需要哪些能力。写报告需要:信息检索能力、数据分析能力、结构化写作能力、事实校验能力。执行层是把每个能力映射到具体的Agent任务。
拆解的时候有个原则:每个Agent任务应该是可独立验证的。也就是说,这个任务做完之后,你能明确判断它做得好不好。如果判断不了,说明拆得不够细,或者任务定义不够明确。
拆完之后画成DAG。我一般用文本先画,确认逻辑没问题了再写成代码。比如上面那个报告任务:
信息检索 ──→ 数据分析 ──→ 结构化写作 ──→ 事实校验 │ │ └──────────→ 趋势判断 ─────────┘这个DAG里,数据分析的结果同时给结构化写作和趋势判断用,事实校验依赖写作和趋势判断两个结果。调度器会自动处理这种依赖关系。
4.2 提示词工程:让每个Agent各司其职
多Agent系统里,提示词的质量比单Agent更重要,因为每个Agent的职责边界必须非常清晰。我写Agent提示词的时候,固定包含五个部分:角色定义、任务描述、输入说明、输出格式、约束条件。
角色定义要具体到"你是一个有10年经验的数据分析师",而不是"你是一个助手"。任务描述要明确"你要做什么"和"你不做什么"。输入说明要告诉Agent它收到的数据是什么格式、从哪来的。输出格式要给出具体的schema或者示例。约束条件要列出禁止事项,比如"不要编造数据""如果信息不足,明确说明而不是猜测"。
我举个例子,一个事实校验Agent的提示词:
你是一个事实核查专家,专门验证文本中的事实性陈述是否准确。 你的任务是:接收一段分析文本,逐条检查其中的事实性陈述,标记出无法验证或与已知信息矛盾的内容。 输入格式:一段包含多个事实性陈述的分析文本。 输出格式:JSON数组,每个元素包含: - statement: 原始陈述 - verdict: "verified" | "unverified" | "contradicted" - reason: 判断理由 - suggestion: 修正建议(如果有) 约束条件: - 只检查事实性陈述,不检查观点和判断 - 如果无法确定,标记为unverified,不要猜测 - 不要修改原文,只输出校验结果这种结构化的提示词,能让Agent的输出稳定性大幅提升。我实测下来,加了输出格式约束之后,解析失败率从15%降到了2%以下。
4.3 结果汇总与冲突消解
多个Agent的输出汇总到一起,经常会出现冲突。比如检索Agent说"市场规模是100亿",分析Agent说"根据数据推算市场规模约120亿"。这种冲突不处理,最终报告就会自相矛盾。
冲突消解我一般分三步:检测、评估、决策。检测就是找出相互矛盾的陈述。评估是判断哪个更可信,依据包括数据来源的权威性、推理过程的严谨性、与其他信息的一致性。决策是选择保留一个、合并两个、还是标记为待确认。
实际实现的时候,我会加一个专门的"仲裁Agent",把冲突双方的信息都给它,让它做判断。仲裁Agent的提示词里会强调"优先采信有明确数据来源的陈述""如果两个陈述都有道理,尝试找出它们成立的条件差异"。
class ArbitrationAgent: async def resolve(self, conflicts: list[dict]) -> list[dict]: prompt = self._build_arbitration_prompt(conflicts) result = await self.model_client.chat(prompt) return self._parse_arbitration(result) def _build_arbitration_prompt(self, conflicts): lines = ["以下陈述存在冲突,请逐条仲裁:\n"] for i, c in enumerate(conflicts): lines.append(f"冲突{i+1}:") lines.append(f" 陈述A: {c['a']}") lines.append(f" 陈述B: {c['b']}") lines.append(f" 背景: {c['context']}\n") lines.append("对每个冲突,输出:保留哪条、理由、或合并方案。") return "\n".join(lines)仲裁Agent不是万能的,有些冲突它也判断不了。这时候我会把冲突标记出来,在最终输出里以"注"的形式呈现,让人类做最终判断。这比强行选一个要好,因为强行选一个可能选错,而标记出来至少不会误导。
4.4 全流程串联与状态管理
把上面所有环节串起来,一个完整的协同任务流程是这样的:
- 主调度器加载DAG配置,初始化所有Agent任务
- 按依赖顺序调度任务,无依赖的任务并行执行
- 每个Agent任务从共享状态读取输入,执行后写回结果
- 所有任务完成后,汇总Agent收集所有结果,做冲突消解
- 最终输出Agent生成交付物,校验Agent做最后检查
状态管理这块,我用一个共享的ContextStore来存所有中间结果。每个Agent任务执行前从里面读,执行后往里写。ContextStore的key用任务名.字段名的格式,避免命名冲突。
class ContextStore: def __init__(self): self._data = {} self._lock = asyncio.Lock() async def get(self, key: str): async with self._lock: return self._data.get(key) async def set(self, key: str, value): async with self._lock: self._data[key] = value async def get_by_prefix(self, prefix: str): async with self._lock: return { k: v for k, v in self._data.items() if k.startswith(prefix) }get_by_prefix这个方法很实用,汇总Agent可以用它一次性拿到某个任务的所有输出,不用一个个key去查。
5. 常见问题与排查技巧实录
5.1 Agent"跑偏"了怎么办
这是最高频的问题。Agent没有按预期执行任务,输出了一堆无关内容。排查思路分三层:
第一层:检查提示词。最常见的原因是提示词里的约束不够明确。比如你写"分析这段数据",Agent可能给你写一篇散文。改成"分析这段数据,输出JSON格式,包含trend、anomaly、summary三个字段",输出就稳定了。
第二层:检查输入。Agent收到的输入里可能包含了干扰信息。比如上游Agent的输出里带了很多解释性文字,当前Agent被这些文字带偏了。解决办法是在_build_context里做输入清洗,只传必要字段。
第三层:检查模型参数。temperature设太高,输出随机性就大。多Agent系统里,除了创意类任务,我一般把temperature设在0.1到0.3之间。top_p设在0.9左右。
下面这张表是我整理的常见"跑偏"现象和对应解法:
| 现象 | 可能原因 | 解法 |
|---|---|---|
| 输出格式不对 | 提示词缺少格式约束 | 加输出schema和示例 |
| 内容偏离主题 | 输入包含干扰信息 | 清洗输入,只传必要字段 |
| 输出过于简略 | 提示词没有长度要求 | 明确要求"至少X字"或"详细说明" |
| 重复上游内容 | 提示词没有区分任务 | 强调"你的任务是X,不是Y" |
| 编造信息 | 缺少事实约束 | 加"不确定就说不确定"的约束 |
5.2 任务卡死或死循环
多Agent系统跑着跑着不动了,一般两个原因:死锁和无限循环。
死锁的典型场景是Agent A等Agent B的结果,Agent B等Agent A的结果。在DAG调度里,死锁表现为_get_ready_tasks永远返回空列表,但还有任务没完成。排查方法是检查DAG里有没有循环依赖。我一般会在调度器初始化的时候做一次拓扑排序,有环就直接报错。
无限循环的典型场景是重试机制没有上限,或者Agent之间的对话没有终止条件。比如两个Agent互相要求对方补充信息,来回几十轮。解决办法是给每个任务设最大重试次数,给Agent对话设最大轮数。
MAX_RETRIES = 3 MAX_DIALOGUE_ROUNDS = 5 async def execute_with_retry(task, max_retries=MAX_RETRIES): for attempt in range(max_retries): try: return await task() except Exception as e: if attempt == max_retries - 1: raise await asyncio.sleep(2 ** attempt)指数退避(2 ** attempt)是重试等待时间的常用策略,第一次等2秒,第二次等4秒,第三次等8秒。这样能避免短时间内大量重试把API打挂。
5.3 输出质量不稳定
同一个任务,跑十次有三次结果很差。这种不稳定问题,根源往往是模型的不确定性和任务定义的模糊性。
降低模型不确定性,除了调temperature,还可以用多次采样投票。让同一个Agent任务跑3次,取多数一致的结果。这个方法对分类、判断类任务特别有效。代价是成本翻3倍,所以只对关键任务用。
降低任务模糊性,核心是给例子。在提示词里放一两个输入输出的示例,Agent的输出稳定性会明显提升。这叫few-shot prompting,在多Agent系统里效果比单Agent更明显,因为每个Agent的任务更聚焦,示例的参考价值更大。
我还有一个私藏的技巧:给Agent加"自检"步骤。让Agent在输出最终结果之前,先自己检查一遍"我的输出是否符合格式要求""我是否完成了所有子任务""我有没有编造信息"。这个自检步骤能让输出合格率提升10到15个百分点。
5.4 成本失控
多Agent系统跑起来,token消耗是单Agent的好几倍。一个复杂任务跑下来,几十万token很正常。成本控制我一般从三个地方入手:
第一,精简上下文。每个Agent只接收必要的输入,不要把上游所有输出都塞进去。我见过一个项目,每个Agent都把完整的历史对话带上,token消耗直接爆炸。改成只带相关字段之后,成本降了60%。
第二,分级模型。不是所有Agent都需要用最强的模型。检索、格式化、简单判断这类任务,用便宜的小模型就够了。只有核心的推理、写作、仲裁任务才用大模型。我一般把任务分成三档:简单任务用小模型,中等任务用中模型,复杂任务用大模型。
第三,缓存。相同的输入不要重复调用模型。我在AgentTask里加了一层缓存,key是提示词的hash,value是模型输出。对于检索类、校验类这种输入重复率高的任务,缓存命中率能到30%以上。
注意:缓存要注意失效策略。如果上游数据变了,缓存必须失效。我一般给缓存加一个TTL(生存时间),比如1小时,过期自动清除。
6. 从单Agent到多Agent的迁移经验
如果你现在有一个跑得还不错的单Agent系统,想迁移到多Agent,我的建议是渐进式迁移,不要推倒重来。
第一步,先把单Agent里的不同职责识别出来。比如一个客服Agent,它其实在做意图识别、知识检索、回复生成三件事。把这三件事拆成三个Agent,用最简单的中心化架构串起来。这一步的目的是验证多Agent的协作流程能不能跑通。
第二步,给每个Agent写独立的提示词,做独立的测试。确保每个Agent在自己的职责范围内表现稳定。这一步最耗时,但最值得。我见过太多人跳过这一步,直接把单Agent的提示词拆成三份就上线,结果每个Agent都不稳定。
第三步,加上调度器和状态管理。把之前手动串的流程改成自动调度。这一步开始引入DAG或者状态机,处理依赖关系和并行执行。
第四步,加上监控和日志。多Agent系统的可观测性比单Agent重要得多。每个Agent的输入、输出、耗时、token消耗都要记录。出问题的时候,你能快速定位是哪个Agent、哪个环节出的问题。
我自己的项目从单Agent迁移到多Agent,前后花了三周。第一周做拆解和提示词,第二周做调度和状态管理,第三周做监控和调优。迁移之后,任务成功率从72%提升到了91%,虽然token成本涨了2.5倍,但考虑到成功率的大幅提升,这个投入是值得的。
最后分享一个我在实际项目中总结的小技巧:给每个Agent起一个有意义的名字。不要用agent_1、agent_2这种,用retriever、analyzer、writer、verifier这种。这个名字会出现在日志里、监控面板上、错误信息中。名字有意义,排查问题的时候能省很多脑力。这个习惯我从第一个多Agent项目保持到现在,每次看日志都觉得当初这个决定太对了。