从零构建子代理系统:基于职责链的智能体协同架构与工程实践
1. 从“单兵作战”到“团队协作”为什么需要子代理在之前的几期内容里我们已经搭建了一个能够独立处理任务的智能体Agent。它能理解指令、调用工具、完成任务就像一个训练有素的“单兵”。但当我们面对一个复杂、多步骤、需要多领域知识协同的任务时这个“单兵”的局限性就暴露出来了。比如一个典型的用户需求是“帮我分析一下上个月的销售数据生成一份PPT报告并总结出三个核心改进点最后把报告发到我的邮箱。”这个任务至少涉及四个核心环节数据获取与清洗、数据分析与洞察、文档生成与美化、邮件发送。让一个Agent去完成所有环节就像要求一个程序员同时精通数据库、统计学、平面设计和网络协议不仅效率低下而且极易出错任何一个环节的失败都会导致整个任务链的崩溃。这就是“子代理”Sub-Agent概念诞生的背景。子代理模式本质上是一种基于职责链Chain of Responsibility和编排Orchestration思想的架构设计。它将一个复杂的宏观任务Macro-Task分解为一系列原子化的微观任务Micro-Tasks并为每个微观任务分派一个专门的、能力聚焦的“专家”Agent去执行。主Agent或称协调者Agent负责任务分解、流程编排、结果汇总和异常处理而子代理们则各司其职专注于自己最擅长的领域。这种架构带来的好处是显而易见的高内聚低耦合每个子代理功能单一内部逻辑清晰易于开发、测试和维护。修改一个数据处理的子代理不会影响到生成PPT的子代理。专业化与高性能可以为不同的子代理配置不同的模型、工具和知识库。例如数据分析子代理可以使用擅长逻辑推理的模型如GPT-4和Pandas库而PPT生成子代理则可以接入擅长结构化输出和审美的模型如Claude 3并调用专门的PPT生成API。鲁棒性增强单个子代理的失败不会导致整个系统瘫痪。主Agent可以设计重试机制、备用方案或者将失败任务路由给其他子代理。可扩展性当需要增加新的能力时我们只需要开发一个新的子代理并将其注册到系统中无需改动主Agent和其他子代理的核心逻辑。理解了“为什么”我们接下来就要深入“怎么做”。实现子代理系统远不止是创建几个Agent实例那么简单它涉及到一套完整的设计范式、通信协议和状态管理机制。2. 核心架构设计主从模式与工作流引擎在动手写代码之前我们必须先确定系统的顶层架构。目前主流的设计模式有两种中心化编排模式和去中心化协同模式。对于绝大多数从零开始的场景我强烈推荐从中心化编排模式入手它结构清晰易于控制和调试。2.1 中心化编排模式详解在这种模式下存在一个绝对的“大脑”——主协调器Master Coordinator。它负责整个任务的生命周期管理。其工作流程可以抽象为以下几个核心阶段任务接收与解析主协调器接收用户输入的原始任务Natural Language Task利用大语言模型LLM的能力将其解析成一个结构化的任务计划Task Plan。这个计划通常是一个有向无环图DAG节点代表子任务边代表依赖关系。例如对于“分析销售数据并做PPT”的任务解析出的DAG可能是[获取销售数据] - [清洗与分析数据] - [生成分析摘要] - [制作PPT幻灯片] - [发送邮件]其中[制作PPT幻灯片]依赖于[生成分析摘要]的输出。子代理调度与执行主协调器根据任务计划依次或并行地调用相应的子代理。它会将当前子任务所需的上下文包括用户原始指令、上游子任务的输出等打包成一个工作指令Work Order发送给对应的子代理。结果收集与状态管理主协调器监听每个子代理的执行状态成功、失败、进行中。它需要维护一个全局的上下文存储器Context Store用于保存所有子任务的输入、输出和中间状态。这是实现任务间数据传递和错误恢复的关键。异常处理与流程控制当某个子代理执行失败或超时主协调器需要根据预设的策略如重试、忽略、更换子代理、整体失败进行处理并决定工作流是继续、回滚还是终止。2.2 关键组件定义与接口设计基于以上流程我们可以定义出几个核心的类MasterAgent(主协调器)属性子代理注册表sub_agent_registry、工作流状态机、上下文存储器。方法register_sub_agent(name, agent_instance): 注册一个子代理。parse_task(user_input) - TaskPlan: 调用LLM进行任务分解与规划。execute_plan(task_plan): 执行任务计划负责调度。_dispatch_to_sub_agent(sub_task, context): 内部方法用于调用具体子代理。SubAgent(子代理基类)这是一个抽象基类ABC所有具体的子代理都需要继承它。属性name唯一标识、description能力描述、required_tools所需工具列表。方法execute(context: Dict) - Dict:必须实现的抽象方法。接收上下文字典执行任务返回结果字典。这是子代理与主协调器之间的契约。TaskPlanSubTask(任务计划与子任务)用于结构化表示分解后的任务。SubTask对象应包含id,type对应子代理类型,goal目标描述,dependencies依赖的父任务ID列表,status,result等字段。ContextStore(上下文存储器)可以是一个简单的内存字典、Redis或数据库。用于存储task_id为键所有子任务输入输出为值的全局状态。2.3 通信协议同步调用 vs. 异步消息队列子代理如何被调用这里有两种选择同步调用主协调器直接调用子代理的execute方法并阻塞等待结果。实现简单适用于轻量级、快速完成的子任务。但主协调器容易被长时间运行的任务阻塞。异步消息队列主协调器将工作指令发布到一个消息队列如RabbitMQ, Redis Stream, Kafka子代理作为消费者从队列中领取任务并执行完成后将结果发布到另一个结果队列。主协调器异步监听结果。这种方式解耦彻底支持高并发和弹性伸缩是生产级系统的首选。对于我们的“从零实现”系列为了聚焦核心逻辑我会先采用同步调用的方式实现一个最小可行系统MVP在后续优化章节再探讨异步改造。3. 手把手实现一个最小可行子代理系统理论说得再多不如一行代码。让我们开始构建。我们将创建以下文件结构my_agent_project/ ├── master_agent.py ├── sub_agent_base.py ├── sub_agents/ │ ├── __init__.py │ ├── data_fetcher_agent.py │ ├── data_analyzer_agent.py │ └── report_generator_agent.py ├── context_store.py ├── task_planner.py └── main.py3.1 第一步定义子代理基类与上下文存储sub_agent_base.pyfrom abc import ABC, abstractmethod from typing import Dict, Any, List class SubAgent(ABC): 所有子代理的抽象基类。 def __init__(self, name: str, description: str): self.name name self.description description self.required_tools: List[str] [] # 例如[web_search, calculator] abstractmethod async def execute(self, context: Dict[str, Any]) - Dict[str, Any]: 执行子任务的核心方法。 Args: context: 来自主协调器的上下文信息包含任务目标、上游结果等。 Returns: 一个字典必须包含 status (success/failure) 和 output 键。 可以包含其他任意键如 artifacts (生成的文件路径)、metrics等。 pass def get_info(self) - Dict[str, str]: 返回子代理的元信息用于主协调器注册和任务规划。 return { name: self.name, description: self.description, required_tools: self.required_tools }context_store.pyfrom typing import Dict, Any, Optional class InMemoryContextStore: 一个简单的内存上下文存储用于演示。生产环境应替换为持久化存储。 def __init__(self): self._store: Dict[str, Dict[str, Any]] {} # task_id - context dict def save_context(self, task_id: str, context: Dict[str, Any]): 保存或更新某个任务的上下文。 self._store[task_id] context def get_context(self, task_id: str) - Optional[Dict[str, Any]]: 获取某个任务的上下文。 return self._store.get(task_id) def update_subtask_result(self, task_id: str, subtask_id: str, result: Dict[str, Any]): 更新某个任务下特定子任务的结果。 if task_id in self._store: if subtask_results not in self._store[task_id]: self._store[task_id][subtask_results] {} self._store[task_id][subtask_results][subtask_id] result3.2 第二步实现几个具体的子代理我们实现三个简单的子代理数据获取、数据分析、报告生成。sub_agents/data_fetcher_agent.pyimport asyncio from sub_agent_base import SubAgent class DataFetcherAgent(SubAgent): 模拟数据获取子代理。 def __init__(self): super().__init__( namedata_fetcher, description从模拟数据源或API获取指定数据。 ) self.required_tools [mock_database] async def execute(self, context: Dict) - Dict: print(f[DataFetcher] 开始执行任务: {context.get(goal)}) # 模拟耗时操作 await asyncio.sleep(1) # 这里应该是真实的数据库或API调用逻辑 # 例如data query_database(context[query]) mock_data [ {month: Jan, sales: 15000, cost: 8000}, {month: Feb, sales: 18000, cost: 9000}, {month: Mar, sales: 22000, cost: 11000}, ] print(f[DataFetcher] 数据获取成功共{len(mock_data)}条记录。) return { status: success, output: mock_data, message: 数据获取完成 }sub_agents/data_analyzer_agent.pyimport asyncio from sub_agent_base import SubAgent class DataAnalyzerAgent(SubAgent): 模拟数据分析子代理。 def __init__(self): super().__init__( namedata_analyzer, description对提供的数据集进行统计分析计算关键指标。 ) self.required_tools [pandas, numpy] # 描述所需能力 async def execute(self, context: Dict) - Dict: print(f[DataAnalyzer] 开始分析数据...) data context.get(input_data, []) if not data: return {status: failure, output: None, error: 输入数据为空} await asyncio.sleep(2) # 模拟计算耗时 # 模拟分析逻辑 total_sales sum(item[sales] for item in data) total_cost sum(item[cost] for item in data) avg_profit_margin ((total_sales - total_cost) / total_sales * 100) if total_sales else 0 best_month max(data, keylambda x: x[sales])[month] analysis_result { total_sales: total_sales, total_cost: total_cost, total_profit: total_sales - total_cost, avg_profit_margin: round(avg_profit_margin, 2), best_performing_month: best_month, summary: f第一季度总销售额为{total_sales}平均利润率{avg_profit_margin}%表现最佳的月份是{best_month}。 } print(f[DataAnalyzer] 分析完成: {analysis_result[summary]}) return { status: success, output: analysis_result }sub_agents/report_generator_agent.pyimport asyncio from sub_agent_base import SubAgent class ReportGeneratorAgent(SubAgent): 模拟报告生成子代理。 def __init__(self): super().__init__( namereport_generator, description根据分析结果生成结构化的文本报告。 ) async def execute(self, context: Dict) - Dict: print(f[ReportGenerator] 正在生成报告...) analysis context.get(input_analysis, {}) await asyncio.sleep(1) report f # 销售数据分析报告 ## 执行摘要 基于{context.get(period, 近期)}的销售数据已完成深度分析。 ## 核心发现 1. **财务总览**总销售额为 {analysis.get(total_sales, N/A)}总成本为 {analysis.get(total_cost, N/A)}实现总利润 **{analysis.get(total_profit, N/A)}**。 2. **盈利能力**平均利润率为 **{analysis.get(avg_profit_margin, N/A)}%**表明业务具有健康的盈利空间。 3. **月度表现**表现最佳的月份是 **{analysis.get(best_performing_month, N/A)}**建议复盘该月份的市场策略与运营动作。 ## 建议 - 继续保持高利润率月份的运营策略。 - 关注成本控制尤其是在销售淡季。 - 深入挖掘最佳月份的成功因素并尝试复制到其他时段。 print(f[ReportGenerator] 报告生成成功长度{len(report)}字符。) return { status: success, output: report, artifacts: { report_text: report, format: markdown } }3.3 第三步实现任务规划器LLM驱动这是系统的“智能”所在。我们利用LLM将自然语言任务分解为结构化计划。这里使用OpenAI API进行演示。task_planner.pyimport openai from typing import List, Dict, Any import json # 假设你已经设置了环境变量 OPENAI_API_KEY # openai.api_key os.getenv(OPENAI_API_KEY) class TaskPlanner: 利用LLM进行任务分解与规划。 def __init__(self, model: str gpt-3.5-turbo): self.model model # 可用的子代理类型描述用于few-shot提示 self.available_agents [ {name: data_fetcher, desc: 从数据库或API获取原始数据。}, {name: data_analyzer, desc: 对数据进行统计、计算和洞察分析。}, {name: report_generator, desc: 根据分析结果生成文本或可视化报告。}, {name: email_sender, desc: 发送邮件或通知。}, ] def plan(self, user_task: str) - List[Dict[str, Any]]: 将用户任务分解为子任务列表。 prompt f 你是一个高级任务规划AI。请将以下用户任务分解为一系列顺序执行的子任务。 可用的子代理类型及其能力如下 {json.dumps(self.available_agents, indent2, ensure_asciiFalse)} 请输出一个JSON数组每个元素是一个子任务对象包含以下字段 - id: 唯一数字标识从1开始。 - agent_type: 必须从上述可用代理的name中选择。 - goal: 对该子任务目标的清晰描述。 - dependencies: 一个数组列出此任务所依赖的父任务id。如果没有依赖则为空数组[]。 用户任务{user_task} 只输出JSON数组不要有任何其他解释。 try: response openai.ChatCompletion.create( modelself.model, messages[{role: user, content: prompt}], temperature0.1, # 低随机性保证规划稳定 ) result response.choices[0].message.content.strip() # 清理可能出现的markdown代码块标记 if result.startswith(json): result result[7:] if result.endswith(): result result[:-3] subtasks json.loads(result) return subtasks except (json.JSONDecodeError, openai.error.OpenAIError) as e: print(f任务规划失败: {e}) # 降级方案返回一个默认的简单计划 return [ {id: 1, agent_type: data_fetcher, goal: 获取销售数据, dependencies: []}, {id: 2, agent_type: data_analyzer, goal: 分析销售数据计算关键指标, dependencies: [1]}, {id: 3, agent_type: report_generator, goal: 根据分析结果生成总结报告, dependencies: [2]}, ]注意在实际项目中LLM的规划结果可能不稳定。生产环境需要更鲁棒的提示工程如Chain-of-Thought并添加结果验证和修正逻辑。这里的降级方案保证了系统在LLM服务异常时仍能执行基本流程。3.4 第四步实现主协调器现在我们把所有组件组装起来。master_agent.pyimport asyncio from typing import Dict, Any, List from task_planner import TaskPlanner from context_store import InMemoryContextStore class MasterAgent: 主协调器负责任务分解、调度与执行。 def __init__(self): self.sub_agents: Dict[str, Any] {} # name - agent_instance self.context_store InMemoryContextStore() self.task_planner TaskPlanner() self.current_task_id None def register_sub_agent(self, agent_instance): 注册一个子代理实例。 info agent_instance.get_info() self.sub_agents[info[name]] agent_instance print(f[Master] 已注册子代理: {info[name]} - {info[description]}) async def execute_task(self, user_input: str) - Dict[str, Any]: 执行用户任务的入口方法。 import uuid self.current_task_id str(uuid.uuid4())[:8] print(f\n 开始执行任务 (ID: {self.current_task_id}) ) print(f用户指令: {user_input}) # 1. 初始化上下文 initial_context { task_id: self.current_task_id, user_input: user_input, subtask_results: {} } self.context_store.save_context(self.current_task_id, initial_context) # 2. 任务规划 print(f\n[Master] 正在规划任务...) subtasks self.task_planner.plan(user_input) print(f[Master] 规划完成共{len(subtasks)}个子任务:) for st in subtasks: print(f - [{st[id]}] {st[agent_type]}: {st[goal]} (依赖: {st[dependencies]})) # 3. 顺序执行子任务简化版未处理并行 for subtask in subtasks: await self._execute_subtask(subtask) # 4. 汇总最终结果 final_context self.context_store.get_context(self.current_task_id) final_output self._compile_final_output(final_context) print(f\n 任务执行完成 (ID: {self.current_task_id}) ) print(f最终输出:\n{final_output[:500]}...) # 只打印前500字符 return final_output async def _execute_subtask(self, subtask: Dict): 执行单个子任务。 subtask_id f{subtask[id]}_{subtask[agent_type]} agent_type subtask[agent_type] if agent_type not in self.sub_agents: print(f[Master] 错误未注册的子代理类型 {agent_type}任务 {subtask_id} 跳过。) self.context_store.update_subtask_result( self.current_task_id, subtask_id, {status: failure, error: fAgent {agent_type} not found} ) return # 检查依赖是否全部成功 for dep_id in subtask[dependencies]: dep_key f{dep_id}_{...} # 需要根据依赖ID找到对应的结果键这里简化处理 # 实际中需要更复杂的依赖检查逻辑 pass # 准备执行上下文 task_context self.context_store.get_context(self.current_task_id) execution_context { goal: subtask[goal], task_id: self.current_task_id, subtask_id: subtask_id, **task_context # 将全局上下文合并进来 } # 调用子代理 print(f\n[Master] 调度子任务 {subtask_id} - {agent_type}) try: agent self.sub_agents[agent_type] result await agent.execute(execution_context) print(f[Master] 子任务 {subtask_id} 完成状态: {result[status]}) # 保存结果到上下文 self.context_store.update_subtask_result(self.current_task_id, subtask_id, result) # 将本次输出作为下一次的潜在输入简化传递逻辑 if result[status] success and output in result: task_context self.context_store.get_context(self.current_task_id) task_context[f{agent_type}_output] result[output] self.context_store.save_context(self.current_task_id, task_context) except Exception as e: print(f[Master] 子任务 {subtask_id} 执行异常: {e}) self.context_store.update_subtask_result( self.current_task_id, subtask_id, {status: failure, error: str(e)} ) def _compile_final_output(self, final_context: Dict) - Dict: 从最终上下文中编译出对用户友好的输出。 # 这里可以设计复杂的汇总逻辑例如找到报告生成器的输出 output { task_id: self.current_task_id, status: completed, summary: 所有子任务已执行完毕。, raw_context: final_context # 实际项目中可能只返回关键结果 } # 尝试提取报告 for key, value in final_context.get(subtask_results, {}).items(): if report_generator in key and value.get(status) success: output[final_report] value.get(output, ) break return output3.5 第五步组装并运行整个系统main.pyimport asyncio from master_agent import MasterAgent from sub_agents.data_fetcher_agent import DataFetcherAgent from sub_agents.data_analyzer_agent import DataAnalyzerAgent from sub_agents.report_generator_agent import ReportGeneratorAgent async def main(): # 1. 初始化主协调器 master MasterAgent() # 2. 注册所有可用的子代理 master.register_sub_agent(DataFetcherAgent()) master.register_sub_agent(DataAnalyzerAgent()) master.register_sub_agent(ReportGeneratorAgent()) # 3. 执行一个示例任务 user_task 帮我分析一下第一季度的销售数据并生成一份分析报告。 final_result await master.execute_task(user_task) # 4. 打印最终报告 if final_report in final_result: print(\n *50) print(生成的最终报告) print(*50) print(final_result[final_report]) if __name__ __main__: asyncio.run(main())运行python main.py你将看到控制台输出完整的任务执行流程从规划、调度到每个子代理的执行最终输出一份结构化的分析报告。这标志着你已经成功实现了一个最小可用的子代理系统4. 从MVP到生产级关键优化与进阶思考上面的MVP演示了核心流程但距离一个健壮、可用的生产系统还有很大距离。以下是几个必须考虑的进阶方向4.1 子任务依赖与并行执行我们的MVP是顺序执行的。真实场景中许多子任务可以并行。我们需要实现一个依赖解析器和任务调度器。依赖解析根据子任务的dependencies字段构建任务DAG。拓扑排序与调度使用拓扑排序算法确定执行顺序将没有依赖关系的任务放入线程池或异步任务队列中并行执行。可以使用asyncio.gather或concurrent.futures。# 伪代码示例并行调度 async def execute_plan_parallel(self, task_plan): # 1. 构建DAG并拓扑排序 sorted_tasks topological_sort(task_plan) # 2. 为每个任务创建异步future task_futures {} for task in sorted_tasks: # 等待所有依赖任务完成 deps [task_futures[dep] for dep in task.dependencies] if deps: await asyncio.gather(*deps) # 执行当前任务 task_futures[task.id] asyncio.create_task(self._execute_subtask(task)) # 3. 等待所有任务完成 await asyncio.gather(*task_futures.values())4.2 上下文管理与数据传递MVP中简单的上下文传递f{agent_type}_output非常脆弱。我们需要一个标准化的数据契约。定义输入/输出规范为每个子代理类型定义明确的输入字段和输出字段。例如DataAnalyzerAgent的输入必须包含input_data字段输出必须包含analysis_result字段。上下文路由主协调器需要根据子任务依赖关系自动将上游输出的特定字段如analysis_result路由到下游任务的输入字段如input_analysis。这可以通过在任务规划时由LLM指定或通过一个配置映射表来实现。4.3 错误处理、重试与熔断错误分类定义系统错误网络超时、业务错误数据为空、逻辑错误参数不合法等。重试策略对于暂时的系统错误如网络抖动配置指数退避重试。可以为每个子代理设置独立的max_retries和retry_delay。熔断机制如果某个子代理连续失败多次暂时将其“熔断”一段时间内不再调度任务给它防止雪崩。可以使用类似circuitbreaker的库。补偿机制对于已经成功但后续任务失败的情况是否需要“回滚”例如如果邮件发送失败是否需要删除已生成的报告文件这需要根据业务定义补偿动作。4.4 引入消息队列进行彻底解耦当系统复杂度和并发量上升时同步调用会成为瓶颈。改造为异步消息模式主协调器将SubTask对象序列化为消息发布到任务队列如agent.tasks。子代理Worker作为独立的进程或服务持续监听任务队列领取任务后执行并将结果发布到结果队列如agent.results。主协调器监听结果队列更新任务状态和上下文。好处子代理可以水平扩展、独立部署、独立升级主协调器无状态更容易实现高可用。4.5 监控、日志与可观测性这是保障系统稳定运行的“眼睛”。结构化日志记录每个任务的完整生命周期事件创建、开始、成功、失败、耗时、输入输出摘要注意脱敏。指标埋点统计每个子代理的调用次数、成功率、平均耗时、排队长度等接入PrometheusGrafana。分布式追踪为每个用户请求生成唯一的trace_id在跨子代理调用时传递便于在日志中串联整个调用链快速定位问题。可以使用OpenTelemetry标准。5. 实战中的经验与避坑指南在真实项目中打磨这套系统我积累了一些宝贵的经验也踩过不少坑经验一子代理的“粒度”设计是艺术也是科学。坑最初我把“数据清洗”和“数据转换”拆成了两个子代理结果发现它们之间需要传递巨大的中间数据通信开销巨大且逻辑耦合依然很紧。解子代理的粒度应该以功能独立性和数据边界来划分。一个子代理最好能完成一个语义上完整、数据输入输出明确、且相对内聚的操作。“获取用户订单数据”是一个好粒度“计算订单金额”可能就太细了它更适合作为“订单分析”子代理内部的一个函数。经验二LLM任务规划的可靠性是最大挑战。坑完全依赖LLM规划它有时会生成不存在的子代理类型或者依赖关系出现循环。解采用“LLM提议 规则校验 人工模板兜底”的三层策略。LLM生成初步计划。用代码校验子代理类型是否已注册依赖图是否有环关键参数是否齐全对于核心、固定的业务流程如“用户注册流程”直接使用预定义的任务模板绕过LLM规划。LLM更适合处理灵活、未知的长尾任务。经验三上下文爆炸与隐私安全。坑将所有历史上下文都传递给下一个子代理导致token数暴涨成本激增且可能泄露敏感信息。解实现上下文摘要与过滤机制。摘要让一个专门的子代理或LLM对冗长的上游结果进行总结只传递摘要。过滤明确定义每个子代理的“上下文需求清单”主协调器只传递清单内的字段。脱敏在上下文存储层或传递前对手机号、邮箱等PII信息进行脱敏处理。经验四测试策略Mock一切可以Mock的。子代理系统涉及多个服务、外部API和LLM调用集成测试极其缓慢且不稳定。策略为每个SubAgent编写独立的单元测试使用Mock对象模拟其依赖如数据库连接、第三方API。为MasterAgent编写集成测试使用Mock子代理来验证任务调度和流程控制逻辑是否正确。使用像pytest-asyncio这样的工具来妥善测试异步代码。针对LLM规划器准备大量涵盖边界的测试用例并验证其输出的结构化数据是否符合预期。实现子代理系统是从构建一个“聪明的单体”迈向构建一个“智能的生态系统”的关键一步。它迫使你以更工程化、更模块化的思维去设计AI应用。虽然初期复杂度有所增加但带来的灵活性、可维护性和扩展性提升对于构建复杂、可靠的AI智能体而言是绝对值得的投入。