ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

多Agent协作架构与任务调度实战:从单Agent到复杂AI协同系统

多Agent协作架构与任务调度实战:从单Agent到复杂AI协同系统 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_factorylist) 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_schemaNone): 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_concurrent5): 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冲突{i1}:) 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_retriesMAX_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是提示词的hashvalue是模型输出。对于检索类、校验类这种输入重复率高的任务缓存命中率能到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项目保持到现在每次看日志都觉得当初这个决定太对了。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进