1. 多Agent协作架构到底在解决什么问题
1.1 从单Agent的瓶颈说起
单Agent跑复杂任务,最典型的翻车场景就是“上下文爆炸”和“能力错配”。你让一个模型同时干需求分析、代码生成、测试验证、文档撰写,它会在中途丢失早期约束,或者把代码风格带进文档里。我实测过一个中等复杂度的后端接口开发任务,单Agent在第三轮迭代时就开始遗忘最初的数据表结构定义,到第五轮直接把字段名写错了。
多Agent协作的核心思路很朴素:把一个大任务拆成多个角色明确的子任务,每个Agent只关注自己那一亩三分地,通过调度层协调它们之间的依赖关系和数据流转。这就像一个小型开发团队,产品经理写需求、后端写接口、测试写用例,各司其职,而不是让一个人从头包到尾。
1.2 协作架构的三种主流形态
目前业界落地的多Agent协作架构,大致可以归为三类:
| 架构类型 | 核心特征 | 适用场景 | 典型代表思路 |
|---|---|---|---|
| 流水线式 | Agent按固定顺序执行,前一个输出是后一个输入 | 流程确定的批处理任务 | 需求→设计→编码→测试 |
| 黑板式 | 所有Agent共享一个状态空间,按需读写 | 需要频繁信息同步的任务 | 联合调试、多源信息融合 |
| 协商式 | Agent之间可以互相提问、反驳、达成共识 | 需要质量校准的复杂决策 | 论文评审、方案论证 |
流水线式最好实现,但灵活性最差;黑板式适合信息密集型任务,但状态管理复杂;协商式质量最高,但Token消耗和延迟也最大。实际项目中,我通常采用混合模式:主流程走流水线,关键节点插入协商环节。
1.3 任务调度的核心挑战
任务调度要解决的不是“谁先谁后”这么简单。真正棘手的是三个问题:
第一,依赖解析。Agent B的输入依赖Agent A的输出,但A的输出可能不完整或格式不对,B需要能识别并触发A的重试或补充。这要求调度层具备输出校验和回退机制。
第二,资源竞争。多个Agent同时调用同一个工具或访问同一份数据时,需要加锁或排队。我见过一个案例,两个Agent同时往同一个文件写内容,结果互相覆盖,整个任务链崩溃。
第三,超时与熔断。某个Agent卡住不返回,不能让它拖死整个流程。需要设置单步超时和全局超时,超时后要么跳过、要么降级、要么触发人工介入。
实操心得:调度层的日志一定要打全。每个Agent的输入、输出、耗时、Token消耗都要记录。出问题时,你才能快速定位是哪个环节的哪个Agent出了什么错。我习惯在调度层加一个
trace_id,贯穿整个任务链,排查效率能提升好几倍。
2. 核心细节解析与实操要点
2.1 Agent角色定义的关键要素
定义一个Agent,不是给它起个名字写句提示词就完事了。一个可落地的Agent定义至少包含以下要素:
- 角色描述:这个Agent是谁,负责什么,不负责什么
- 输入规范:它接收什么格式的数据,必填字段有哪些
- 输出规范:它必须返回什么格式的数据,字段含义是什么
- 可用工具:它能调用哪些外部工具或API
- 约束条件:它不能做什么,比如不能修改某些文件、不能调用某些接口
- 失败处理:出错时返回什么,是否重试,重试几次
我踩过的一个坑是:早期定义Agent时只写了角色描述,没写输出规范。结果代码生成Agent返回的代码块格式不固定,有时用```python包裹,有时直接裸写,导致下游的测试Agent解析失败。后来强制要求所有Agent的输出必须是结构化JSON,问题才解决。
2.2 任务调度的实现方式
任务调度有两种主流实现路径:
路径一:基于代码的硬编码调度。用Python或TypeScript写一个调度器,显式定义每个步骤的执行顺序和条件分支。优点是可控性强,调试方便;缺点是灵活性差,任务流程一变就要改代码。
路径二:基于配置的声明式调度。用YAML或JSON定义任务流,调度器解析配置后动态执行。优点是灵活,改流程不用改代码;缺点是调试相对麻烦,配置写错了不容易发现。
我的建议是:原型阶段用硬编码,快速验证可行性;生产阶段用声明式,方便迭代和维护。下面是一个声明式任务流的配置示例:
task_flow: name: "api_development" steps: - id: "requirement_analysis" agent: "product_manager" input: "${user_input}" output_key: "requirements" timeout: 120 - id: "schema_design" agent: "architect" input: "${requirements}" output_key: "schema" depends_on: ["requirement_analysis"] timeout: 180 - id: "code_generation" agent: "developer" input: requirements: "${requirements}" schema: "${schema}" output_key: "code" depends_on: ["schema_design"] timeout: 300 retry: 2 - id: "test_generation" agent: "tester" input: "${code}" output_key: "tests" depends_on: ["code_generation"] timeout: 180这个配置里,depends_on定义了依赖关系,output_key定义了输出存储的变量名,retry定义了失败重试次数。调度器按拓扑排序依次执行,遇到依赖未满足的步骤就等待。
2.3 通信机制的选择
Agent之间的通信方式直接影响系统的复杂度和性能。常见的有三种:
共享内存/状态。所有Agent读写同一个状态对象。实现简单,但并发写入时需要加锁,且状态膨胀后性能下降明显。
消息队列。Agent之间通过消息队列异步通信。解耦性好,适合分布式部署,但引入额外的基础设施依赖,调试链路变长。
直接调用。一个Agent直接调用另一个Agent的函数。最简单直接,但耦合度高,不适合Agent数量多的场景。
我个人的经验是:Agent数量少于5个时用共享状态,5到15个时用消息队列,超过15个考虑分层调度。分层调度就是设置一个主调度Agent,它下面管几个子调度Agent,每个子调度Agent管一组功能相近的Agent。
2.4 上下文传递的注意事项
多Agent协作中,上下文传递是最容易出问题的地方。每个Agent的上下文窗口有限,不可能把前面所有Agent的完整输出都塞进去。需要做上下文压缩和摘要。
具体做法是:每个Agent执行完后,除了返回完整输出,还要返回一个摘要版本。调度层只把摘要版本传给下游Agent,完整版本存到外部存储(如文件或数据库),需要时再按需读取。
注意:摘要的质量直接影响下游Agent的表现。摘要太简略会丢失关键信息,太详细又起不到压缩作用。我通常要求摘要控制在200字以内,且必须包含“关键决策、关键数据、待解决问题”三个要素。
3. 实操过程与核心环节实现
3.1 环境准备与基础框架搭建
先明确技术栈。我用的方案是Python + FastAPI做调度服务,Redis做状态存储,SQLite做日志持久化。模型侧通过统一的API网关调用,不直接在Agent代码里写模型调用逻辑,方便后续切换模型。
# 创建项目结构 mkdir multi-agent-system && cd multi-agent-system mkdir -p agents scheduler storage logs # 安装核心依赖 pip install fastapi uvicorn redis pydantic httpx项目结构说明:
agents/:存放各个Agent的定义和实现scheduler/:调度器核心逻辑storage/:状态存储和日志持久化logs/:运行日志
3.2 Agent基类的设计与实现
所有Agent继承同一个基类,保证接口统一。基类负责处理输入校验、输出格式化、超时控制、日志记录等通用逻辑。
import json import time import logging from abc import ABC, abstractmethod from typing import Any, Dict, Optional logger = logging.getLogger(__name__) class BaseAgent(ABC): def __init__(self, name: str, timeout: int = 120): self.name = name self.timeout = timeout @abstractmethod def execute(self, input_data: Dict[str, Any]) -> Dict[str, Any]: """子类实现具体的执行逻辑""" pass def validate_input(self, input_data: Dict[str, Any]) -> bool: """输入校验,子类可重写""" return True def format_output(self, raw_output: Any) -> Dict[str, Any]: """输出格式化,子类可重写""" return {"result": raw_output} def run(self, input_data: Dict[str, Any]) -> Dict[str, Any]: """统一的执行入口""" start_time = time.time() trace_id = input_data.get("trace_id", "unknown") logger.info(f"[{trace_id}] Agent {self.name} started") if not self.validate_input(input_data): return { "status": "error", "error": "input validation failed", "agent": self.name } try: raw_output = self.execute(input_data) formatted = self.format_output(raw_output) elapsed = time.time() - start_time logger.info(f"[{trace_id}] Agent {self.name} finished in {elapsed:.2f}s") return { "status": "success", "data": formatted, "agent": self.name, "elapsed": elapsed } except Exception as e: logger.error(f"[{trace_id}] Agent {self.name} failed: {str(e)}") return { "status": "error", "error": str(e), "agent": self.name }这个基类做了几件事:统一日志格式(带trace_id)、统一错误处理、统一输出结构。子类只需要实现execute方法,不用关心日志和错误处理。
3.3 调度器的核心逻辑
调度器负责解析任务流配置、管理依赖关系、按序执行Agent、处理失败重试。
import json from typing import Any, Dict, List from collections import defaultdict class TaskScheduler: def __init__(self, agents: Dict[str, Any], max_retry: int = 2): self.agents = agents self.max_retry = max_retry self.context = {} def resolve_dependencies(self, steps: List[Dict]) -> List[List[Dict]]: """拓扑排序,返回可并行执行的批次""" graph = defaultdict(list) in_degree = defaultdict(int) step_map = {s["id"]: s for s in steps} for step in steps: deps = step.get("depends_on", []) in_degree[step["id"]] = len(deps) for dep in deps: graph[dep].append(step["id"]) batches = [] queue = [sid for sid, deg in in_degree.items() if deg == 0] while queue: batches.append([step_map[sid] for sid in queue]) next_queue = [] for sid in queue: for neighbor in graph[sid]: in_degree[neighbor] -= 1 if in_degree[neighbor] == 0: next_queue.append(neighbor) queue = next_queue return batches def execute_step(self, step: Dict) -> Dict: """执行单个步骤,含重试逻辑""" agent_name = step["agent"] agent = self.agents.get(agent_name) if not agent: return {"status": "error", "error": f"agent {agent_name} not found"} input_data = self._resolve_input(step.get("input", {})) input_data["trace_id"] = self.context.get("trace_id", "unknown") for attempt in range(self.max_retry + 1): result = agent.run(input_data) if result["status"] == "success": return result if attempt < self.max_retry: print(f"Step {step['id']} failed, retrying ({attempt+1}/{self.max_retry})") return result def _resolve_input(self, input_spec: Any) -> Any: """解析输入中的变量引用,如 ${requirements}""" if isinstance(input_spec, str) and input_spec.startswith("${"): key = input_spec[2:-1] return self.context.get(key, {}) elif isinstance(input_spec, dict): return {k: self._resolve_input(v) for k, v in input_spec.items()} return input_spec def run(self, task_flow: Dict, initial_input: Dict) -> Dict: """执行整个任务流""" self.context = {"trace_id": initial_input.get("trace_id", "unknown")} self.context["user_input"] = initial_input steps = task_flow["steps"] batches = self.resolve_dependencies(steps) for batch in batches: for step in batch: result = self.execute_step(step) if result["status"] == "error": return { "status": "error", "failed_step": step["id"], "error": result.get("error") } output_key = step.get("output_key") if output_key: self.context[output_key] = result["data"] return {"status": "success", "context": self.context}这个调度器实现了拓扑排序、变量解析、失败重试三个核心功能。resolve_dependencies方法把任务流拆成可并行执行的批次,_resolve_input方法处理${variable}形式的变量引用,execute_step方法负责单步执行和重试。
3.4 一个完整的协作案例:API开发任务
假设我们要完成一个“用户注册接口”的开发任务,涉及四个Agent:需求分析Agent、架构设计Agent、代码生成Agent、测试生成Agent。
# 定义各个Agent class RequirementAgent(BaseAgent): def execute(self, input_data): user_input = input_data.get("user_input", {}) # 实际项目中这里调用大模型API return { "summary": "用户注册接口,支持邮箱和手机号注册", "fields": ["email", "phone", "password", "nickname"], "constraints": ["密码至少8位", "邮箱格式校验", "手机号格式校验"] } class ArchitectAgent(BaseAgent): def execute(self, input_data): requirements = input_data.get("requirements", {}) return { "table": "users", "columns": [ {"name": "id", "type": "bigint", "primary": True}, {"name": "email", "type": "varchar(255)", "unique": True}, {"name": "phone", "type": "varchar(20)", "unique": True}, {"name": "password_hash", "type": "varchar(255)"}, {"name": "nickname", "type": "varchar(50)"}, {"name": "created_at", "type": "timestamp"} ], "api_path": "/api/v1/user/register", "method": "POST" } class DeveloperAgent(BaseAgent): def execute(self, input_data): schema = input_data.get("schema", {}) # 实际项目中这里调用大模型生成代码 return { "language": "python", "framework": "fastapi", "code": "# 生成的代码..." } class TesterAgent(BaseAgent): def execute(self, input_data): code = input_data.get("code", {}) return { "test_cases": [ {"name": "正常注册", "input": {"email": "a@b.com", "phone": "13800138000", "password": "12345678"}, "expected": 200}, {"name": "密码过短", "input": {"email": "a@b.com", "phone": "13800138000", "password": "123"}, "expected": 400}, {"name": "邮箱格式错误", "input": {"email": "invalid", "phone": "13800138000", "password": "12345678"}, "expected": 400} ] } # 组装并运行 agents = { "product_manager": RequirementAgent("product_manager"), "architect": ArchitectAgent("architect"), "developer": DeveloperAgent("developer"), "tester": TesterAgent("tester") } task_flow = { "name": "api_development", "steps": [ {"id": "req", "agent": "product_manager", "input": "${user_input}", "output_key": "requirements"}, {"id": "design", "agent": "architect", "input": "${requirements}", "output_key": "schema", "depends_on": ["req"]}, {"id": "code", "agent": "developer", "input": {"requirements": "${requirements}", "schema": "${schema}"}, "output_key": "code", "depends_on": ["design"]}, {"id": "test", "agent": "tester", "input": "${code}", "output_key": "tests", "depends_on": ["code"]} ] } scheduler = TaskScheduler(agents) result = scheduler.run(task_flow, {"user_input": "开发一个用户注册接口"}) print(json.dumps(result, indent=2, ensure_ascii=False))这个案例展示了完整的多Agent协作流程:需求分析→架构设计→代码生成→测试生成。每个Agent只关注自己的输入和输出,调度器负责串联。
3.5 并行执行的优化
上面的调度器是串行执行的,实际上没有依赖关系的步骤可以并行。比如“代码生成”和“文档撰写”可以同时进行。改造run方法支持并行:
import concurrent.futures def run_parallel(self, task_flow: Dict, initial_input: Dict) -> Dict: self.context = {"trace_id": initial_input.get("trace_id", "unknown")} self.context["user_input"] = initial_input steps = task_flow["steps"] batches = self.resolve_dependencies(steps) for batch in batches: with concurrent.futures.ThreadPoolExecutor(max_workers=len(batch)) as executor: futures = {executor.submit(self.execute_step, step): step for step in batch} for future in concurrent.futures.as_completed(futures): step = futures[future] result = future.result() if result["status"] == "error": return {"status": "error", "failed_step": step["id"], "error": result.get("error")} output_key = step.get("output_key") if output_key: self.context[output_key] = result["data"] return {"status": "success", "context": self.context}并行执行能把总耗时从“各步骤之和”降到“最长路径耗时”。实测下来,四步串行任务改成两步并行后,总耗时从45秒降到了28秒,提升接近40%。
4. 常见问题与排查技巧实录
4.1 Agent输出格式不一致
这是最高频的问题。同一个Agent在不同轮次可能返回不同格式,比如有时返回纯文本,有时返回JSON,有时JSON外面还包了一层Markdown代码块。
排查思路:先看Agent的提示词是否明确要求了输出格式。如果提示词里写了“返回JSON”,但模型仍然返回Markdown包裹的JSON,需要在format_output方法里做清洗。
解决方案:在基类的format_output里加一个通用的JSON提取逻辑:
import re def extract_json(text: str) -> dict: """从文本中提取JSON,兼容Markdown代码块包裹的情况""" # 尝试直接解析 try: return json.loads(text) except json.JSONDecodeError: pass # 尝试提取```json ... ```中的内容 pattern = r'```(?:json)?\s*\n?(.*?)\n?```' matches = re.findall(pattern, text, re.DOTALL) for match in matches: try: return json.loads(match.strip()) except json.JSONDecodeError: continue # 尝试提取第一个{到最后一个}之间的内容 start = text.find('{') end = text.rfind('}') if start != -1 and end != -1 and end > start: try: return json.loads(text[start:end+1]) except json.JSONDecodeError: pass raise ValueError(f"无法从输出中提取JSON: {text[:200]}")4.2 上下文丢失导致下游Agent报错
下游Agent需要的字段在上游输出中不存在,或者字段名对不上。比如上游返回{"user_name": "张三"},下游期望的是{"username": "张三"}。
排查思路:在调度器的_resolve_input方法里加字段校验,如果引用的变量不存在或字段缺失,立即报错并打印当前上下文。
解决方案:定义Agent时强制要求输出规范,并在调度层做字段映射。我通常会在配置里加一个field_mapping,把上游输出字段映射到下游期望的字段名。
4.3 某个Agent执行超时
某个Agent卡住不返回,整个任务链阻塞。
排查思路:先看是模型调用超时还是Agent内部逻辑死循环。如果是模型调用超时,检查网络和API限流;如果是内部逻辑,检查是否有未处理的异常导致重试无限循环。
解决方案:在Agent基类的run方法里加超时控制:
import signal class TimeoutError(Exception): pass def timeout_handler(signum, frame): raise TimeoutError("Agent execution timeout") class BaseAgent(ABC): def run(self, input_data): signal.signal(signal.SIGALRM, timeout_handler) signal.alarm(self.timeout) try: result = self._run_internal(input_data) signal.alarm(0) return result except TimeoutError: return {"status": "error", "error": "timeout", "agent": self.name}注意:
signal.alarm只在Unix系统有效,Windows下需要用threading.Timer替代。另外,超时后要确保资源被正确释放,避免僵尸线程。
4.4 常见问题速查表
| 问题现象 | 可能原因 | 排查方法 | 解决方案 |
|---|---|---|---|
| Agent输出格式不一致 | 提示词不明确或模型随机性 | 打印原始输出对比 | 加JSON提取清洗逻辑 |
| 下游Agent报字段缺失 | 上游输出字段名不匹配 | 检查上下文中的实际字段 | 加字段映射配置 |
| 任务链卡住不推进 | 某Agent超时或死循环 | 看日志最后一条记录 | 加超时控制和熔断 |
| Token消耗异常高 | 上下文未压缩或重复传递 | 统计每步Token用量 | 加摘要压缩和按需读取 |
| 并行步骤结果互相覆盖 | 共享状态未加锁 | 检查并发写入的key | 加锁或改用消息队列 |
| 重试后仍然失败 | 错误是确定性的而非偶发 | 看错误信息是否相同 | 区分可重试和不可重试错误 |
4.5 独家避坑技巧
技巧一:给每个Agent的输出加版本号。当Agent的提示词或逻辑变更时,输出格式可能变化。加版本号后,下游Agent可以根据版本号做兼容处理。
技巧二:调度层加“干跑”模式。不实际调用模型,而是用Mock数据走一遍流程,验证调度逻辑和字段映射是否正确。这在调试阶段能省大量时间和Token。
技巧三:关键步骤加人工确认节点。对于高风险操作(如删除数据、发送邮件),在调度流中插入一个human_approval步骤,等待人工确认后再继续。实现方式可以是轮询数据库或监听消息队列。
技巧四:日志按trace_id分文件存储。一个任务链的日志写到一个文件里,排查时直接打开对应文件,不用在混合日志里搜索。我通常按logs/{date}/{trace_id}.log的路径存储。
技巧五:定期清理过期上下文。长时间运行的系统,上下文会越积越多。设置一个TTL,超过一定时间的上下文自动清理,避免内存泄漏。
5. 协作架构的扩展与优化方向
5.1 引入协商机制提升输出质量
流水线式协作的问题是:上游Agent的错误会一路传递到下游,没人纠正。引入协商机制后,下游Agent可以对上游输出提出质疑,触发上游重新生成。
实现方式是在调度器中加一个review环节。比如代码生成Agent输出代码后,测试Agent先做一轮静态检查,如果发现明显问题(如语法错误、缺少必要字段),直接返回need_revision状态,调度器触发代码生成Agent重试。
def execute_with_review(self, step, reviewer_agent): """执行步骤后由reviewer审核,不通过则重试""" for attempt in range(self.max_retry + 1): result = self.execute_step(step) if result["status"] == "error": continue review = reviewer_agent.run({"content": result["data"]}) if review["status"] == "success" and review["data"].get("approved"): return result # 审核不通过,把审核意见反馈给原Agent step["input"]["review_feedback"] = review["data"].get("feedback", "") return result5.2 动态任务分解
固定任务流适合流程确定的任务,但面对开放式任务(如“帮我写一篇论文”),需要动态分解。做法是加一个plannerAgent,它根据用户输入动态生成任务流配置,然后调度器按生成的配置执行。
class PlannerAgent(BaseAgent): def execute(self, input_data): user_input = input_data.get("user_input", "") # 调用大模型生成任务流配置 task_flow = { "name": "dynamic_flow", "steps": [ {"id": "research", "agent": "researcher", "input": "${user_input}", "output_key": "research_data"}, {"id": "outline", "agent": "outliner", "input": "${research_data}", "output_key": "outline", "depends_on": ["research"]}, {"id": "writing", "agent": "writer", "input": "${outline}", "output_key": "draft", "depends_on": ["outline"]}, {"id": "review", "agent": "reviewer", "input": "${draft}", "output_key": "final", "depends_on": ["writing"]} ] } return task_flow这种方式的灵活性最高,但对Planner Agent的能力要求也最高。Planner生成的配置如果格式错误或依赖关系有环,调度器需要能检测并报错。
5.3 性能优化的几个实操方向
方向一:缓存重复计算结果。如果多个任务链中有相同的子步骤(如“查询数据库表结构”),可以把结果缓存起来,避免重复调用模型。
方向二:小模型做路由,大模型做生成。调度决策、格式校验、简单分类等任务用小模型(如7B级别)处理,复杂生成任务用大模型。这样能显著降低成本和延迟。
方向三:流式输出与增量处理。对于长文本生成任务,上游Agent流式输出,下游Agent增量处理,不用等上游完全生成完再开始。这需要调度器支持流式传递。
方向四:批处理相似任务。多个用户请求如果涉及相同的Agent和相似的输入,可以合并成一个批次处理,减少模型调用次数。
5.4 监控与可观测性建设
多Agent系统上线后,必须有一套监控体系。我通常关注以下指标:
- 任务成功率:成功完成的任务链占比
- 平均耗时:每个任务链从开始到结束的平均时间
- 各Agent耗时分布:哪个Agent是瓶颈
- Token消耗:每个任务链的Token用量和成本
- 重试率:各Agent的重试次数占比
- 错误分布:各类错误的发生频率
这些指标可以通过调度层埋点收集,写入时序数据库(如Prometheus),再用Grafana做可视化。没有监控的多Agent系统,出了问题就是盲人摸象。
实操心得:监控告警的阈值不要设得太敏感。我一开始把“单步耗时超过60秒”设为告警,结果因为模型API的正常波动,每天收到几十条误报。后来改成“连续3次超过120秒”才告警,噪音少了很多。
6. 从零搭建一个可运行的多Agent系统
6.1 最小可行系统的搭建步骤
如果你现在就想动手搭一个,按以下步骤走:
第一步:定义Agent基类和调度器。直接用前面给的代码,复制到项目里。
第二步:实现两个最简单的Agent。一个EchoAgent(原样返回输入),一个UpperAgent(把输入转大写)。用它们验证调度器是否正常工作。
第三步:接入真实模型。把EchoAgent替换成调用大模型API的Agent。建议先用一个简单的提示词,比如“请把以下内容翻译成英文”,验证模型调用链路。
第四步:定义任务流配置。写一个包含3到4个步骤的YAML配置,跑通完整流程。
第五步:加日志和监控。在调度器的关键节点加日志,确保每个步骤的输入输出都有记录。
第六步:加错误处理和重试。模拟Agent失败的情况,验证重试逻辑是否生效。
这六步走完,你就有了一个可运行的多Agent系统原型。后续的优化和扩展都基于这个原型进行。
6.2 模型调用的统一封装
不要在Agent代码里直接写模型调用逻辑。统一封装成一个ModelClient类,方便切换模型和加缓存。
import httpx import hashlib import json class ModelClient: def __init__(self, base_url: str, api_key: str, model: str): self.base_url = base_url self.api_key = api_key self.model = model self.cache = {} def chat(self, messages: list, temperature: float = 0.7) -> str: # 缓存key基于messages和temperature生成 cache_key = hashlib.md5( json.dumps({"messages": messages, "temp": temperature}, sort_keys=True).encode() ).hexdigest() if cache_key in self.cache: return self.cache[cache_key] response = httpx.post( f"{self.base_url}/chat/completions", headers={"Authorization": f"Bearer {self.api_key}"}, json={ "model": self.model, "messages": messages, "temperature": temperature }, timeout=60 ) result = response.json()["choices"][0]["message"]["content"] self.cache[cache_key] = result return result这个封装做了三件事:统一调用接口、加缓存、统一超时。缓存对于调试阶段特别有用,同样的输入不用重复调用模型,省时省Token。
6.3 配置管理与环境隔离
不同环境(开发、测试、生产)的配置不同。用环境变量或配置文件管理,不要硬编码。
import os from dataclasses import dataclass @dataclass class Config: model_base_url: str model_api_key: str model_name: str redis_host: str redis_port: int max_retry: int default_timeout: int @classmethod def from_env(cls): return cls( model_base_url=os.getenv("MODEL_BASE_URL", "http://localhost:8000/v1"), model_api_key=os.getenv("MODEL_API_KEY", ""), model_name=os.getenv("MODEL_NAME", "default"), redis_host=os.getenv("REDIS_HOST", "localhost"), redis_port=int(os.getenv("REDIS_PORT", "6379")), max_retry=int(os.getenv("MAX_RETRY", "2")), default_timeout=int(os.getenv("DEFAULT_TIMEOUT", "120")) )开发环境可以用本地模型或Mock,生产环境用正式API。环境隔离能避免调试时的误操作影响生产数据。
6.4 部署与扩展的注意事项
多Agent系统的部署有两种模式:单体部署和分布式部署。
单体部署把所有Agent和调度器放在一个进程里,实现简单,适合Agent数量少、任务量不大的场景。分布式部署把Agent拆成独立的服务,通过消息队列通信,适合Agent数量多、需要独立扩缩容的场景。
我建议从单体开始,遇到性能瓶颈再拆。过早分布式化会引入大量运维复杂度,得不偿失。单体部署时,用多线程或异步IO处理并发请求就够了。
如果确实需要分布式,优先拆调度器和Agent。调度器保持单点(或主备),Agent按功能分组部署多个实例。Agent之间不直接通信,所有协调都通过调度器。
注意:分布式部署后,日志追踪变得更复杂。确保trace_id在所有服务间正确传递,否则排查问题时会非常痛苦。我通常用OpenTelemetry做分布式追踪,能自动串联跨服务的调用链。
6.5 一个完整的项目目录结构参考
multi-agent-system/ ├── agents/ │ ├── __init__.py │ ├── base.py # Agent基类 │ ├── requirement.py # 需求分析Agent │ ├── architect.py # 架构设计Agent │ ├── developer.py # 代码生成Agent │ └── tester.py # 测试生成Agent ├── scheduler/ │ ├── __init__.py │ ├── core.py # 调度器核心 │ └── dependency.py # 依赖解析 ├── storage/ │ ├── __init__.py │ ├── context.py # 上下文存储 │ └── logger.py # 日志持久化 ├── config/ │ ├── __init__.py │ └── settings.py # 配置管理 ├── flows/ │ └── api_development.yaml # 任务流配置 ├── tests/ │ ├── test_scheduler.py │ └── test_agents.py ├── main.py # 入口 └── requirements.txt这个结构清晰地区分了Agent定义、调度逻辑、存储、配置和测试。新加一个Agent只需要在agents/下新建文件,在配置里注册即可。
6.6 测试策略
多Agent系统的测试比单Agent复杂,因为涉及多个组件的交互。我通常分三层测试:
单元测试:单独测试每个Agent的execute方法,用Mock输入验证输出格式和内容。
集成测试:测试调度器和Agent的交互,验证任务流能正确执行、变量能正确传递、错误能正确传播。
端到端测试:用真实模型跑完整任务流,验证最终输出是否符合预期。这层测试成本高,不需要每次提交都跑,可以在发版前跑一次。
# 单元测试示例 def test_requirement_agent(): agent = RequirementAgent("test") result = agent.run({"user_input": "开发用户注册接口"}) assert result["status"] == "success" assert "fields" in result["data"] assert "email" in result["data"]["fields"] # 集成测试示例 def test_scheduler_with_mock_agents(): agents = { "a": MockAgent("a", output={"value": 1}), "b": MockAgent("b", output={"value": 2}) } flow = { "steps": [ {"id": "step_a", "agent": "a", "output_key": "a_out"}, {"id": "step_b", "agent": "b", "input": "${a_out}", "output_key": "b_out", "depends_on": ["step_a"]} ] } scheduler = TaskScheduler(agents) result = scheduler.run(flow, {}) assert result["status"] == "success" assert result["context"]["b_out"]["value"] == 2测试覆盖率达到70%以上,基本能保证系统的稳定性。重点测试边界情况:空输入、超长输入、格式错误的输入、Agent超时、Agent返回错误等。
6.7 持续迭代的节奏把控
多Agent系统不是一次性能搭好的。我的迭代节奏是:
第一周:搭最小可行系统,跑通两个Agent的串行流程。
第二周:加错误处理、重试、日志,完善基础设施。
第三周:加并行执行、缓存、监控,提升性能和可观测性。
第四周:引入协商机制和动态任务分解,提升输出质量。
之后就是根据实际使用中的反馈持续优化。每次优化只改一个点,改完跑一轮回归测试,确保没有引入新问题。多Agent系统的复杂度决定了“小步快跑”比“大重构”更安全。
我在实际项目中最深的体会是:多Agent系统的价值不在于Agent数量多,而在于每个Agent的职责是否清晰、调度逻辑是否健壮、错误处理是否完善。一个设计良好的三Agent系统,比一个混乱的十Agent系统产出质量高得多。先把基础架构搭稳,再逐步增加Agent和功能,这条路走下来最踏实。