在构建多智能体协作系统时agency-agents这类模块通常是整个技术方案的核心。它不只是命名习惯而是承载了智能体注册、任务调度、消息传递和结果回收等关键职责。很多人在接触这个概念时会把注意力放在单个 Agent 的提示词或函数实现上结果一旦出现多个 Agent 需要协作代码就会变成一团调用链。下面从agency-agents的工程视角出发说明一套多智能体协作项目应该如何设计、实现、验证和排查。如果你正在设计自己的多智能体项目或者需要接手一个名为 agency-agents 的模块这里会给出一条可直接落地的参考路径。1. 先理解 agency-agents 要解决的核心问题1.1 多智能体协作为什么比单 Agent 困难单个 Agent 适合完成“输入到输出”比较直接的任务给它一段文本它返回摘要给它一个查询它返回答案。但真实业务往往没有这么简单。以“客服工单自动分类并生成处理建议”为例单个 Agent 需要同时完成去重、敏感信息识别、分类、标签提取、建议生成等多个动作。如果全部写在一个函数里提示词会越来越长异常分支越来越多最终难以维护和测试。agency-agents这种设计想解决的问题就是把一个复杂任务拆成多个职责明确的小任务由不同 Agent 分别执行再通过调度器把它们串联起来。这里的 “agency” 更像是“服务方”或“组织”负责管理 Agent 列表、任务路由和协作规则agents是真正干活的具体执行者。两者分离之后新增一个能力只需要注册一个新的 Agent不需要改动其他 Agent 的内部逻辑。维度单个 Agent 大而全agency-agents 多 Agent 协作普通任务队列任务拆分不明显按职责拆成多个阶段按键或类型拆成多个队列节点能力一个函数完成全部每个 Agent 只负责一个子任务每个 Worker 处理一类消息路由方式无按 pipeline 或任务类型路由按队列路由失败处理整体失败可定位到具体 Agent消息重试或进死信队列调试成本低但扩展性差需要链路追踪需要队列监控所以多智能体协作的本质不是“多个 Agent 随便调来调去”而是先定义清楚任务生命周期再定义每个 Agent 的输入、输出和边界。1.2 一个典型任务如何被拆解以一条用户反馈为例可以拆成下面几个子任务内容清洗去掉多余空格、URL 和重复字符。意图识别判断这条反馈是建议、投诉还是咨询。信息补全从上下文中提取用户名、订单号等关键字段。策略生成根据意图和历史记录生成处理建议。这四步分别由四个 Agent 完成每个 Agent 拥有独立的输入输出结构。agency-agents的调度器只需要知道一个顺序先清洗再识别再补全最后生成策略。这个顺序本身是数据可以放到配置文件中而不是写死在代码里。这个设计非常重要因为一旦业务顺序变化只需调整配置不需要重新编译和发布。1.3 agency 与 agents 的职责边界在工程实现中建议把职责分得非常清楚Agency 层维护 Agent 注册表、读取任务编排配置、创建任务、分发任务、汇总结果。Agent 层只处理一件事情的业务逻辑不关心上一个 Agent 从哪来、下一个 Agent 到哪去。消息层定义任务和结果的数据结构保证每个 Agent 的输入输出可序列化、可追踪。调度层控制并发数、超时、重试和路由策略。如果 Agent 层开始依赖其他 Agent 的实现或者 A 的返回值里写死了给 B 看的字段那么协作逻辑就散落在各处后续很难维护。正确做法是Agent 只返回结构化数据调度器把数据作为下一个 Agent 的输入。这样每个 Agent 都可以单独测试不依赖完整链路。1.4 与普通任务队列的差异虽然它们都有生产消费模型但普通任务队列更关注“消息不丢、不重复、可靠消费”而agency-agents更关注“一个任务经过多步处理后每一步都产生了什么中间结果以及最终如何汇总”。队列可以承载 Agent 之间的消息但队列本身不解释业务语义。多智能体项目要把任务类型、上下文、历史记录和结果状态一起放进消息结构里而不只是放一条原始字符串。2. 设计一个最小可运行的多智能体协作项目2.1 环境准备和依赖这里用一个不依赖第三方框架的 Python 最小示例说明agency-agents的骨架。你可以把它换成 Java、Go 或 Node.js但核心结构和数据流是一样的。推荐环境Python 3.10 及以上。只需要标准库threading、concurrent.futures、dataclasses。生产项目可以补充 Pydantic 做数据校验但示例先不引入避免复杂化。如果你在正式项目中使用建议先确认 Python 版本并且在虚拟环境中安装依赖。下面是环境检查命令python --version pip --version如果在 Windows 上使用多线程调度注意放在if __name__ __main__中启动避免子进程递归问题。2.2 项目目录结构一个可维护的多智能体项目目录结构可以从一开始就设计好。下面是一个参考结构agency-agents/ ├── agency/ │ ├── __init__.py │ ├── models.py # Task、TaskResult 等数据结构 │ ├── registry.py # Agent 注册中心 │ ├── dispatcher.py # 调度器 │ └── base.py # Agent 抽象基类 ├── agents/ │ ├── __init__.py │ ├── clean_agent.py # 实际 Agent 实现 │ └── summary_agent.py ├── config/ │ └── pipeline.yaml # 任务编排配置 ├── main.py # 程序入口 └── tests/ └── test_pipeline.py这个结构的好处是agency目录可以复用agents目录按业务扩展config目录让编排规则外部化tests目录保证核心链路可回归。2.3 数据模型定义先定义任务。任务要能唯一标识并包含业务数据、任务类型和历史记录# agency/models.py from dataclasses import dataclass, field from typing import Any, Dict, List import uuid import time dataclass class Task: task_id: str field(default_factorylambda: uuid.uuid4().hex) task_type: str payload: Dict[str, Any] field(default_factorydict) history: List[Dict[str, Any]] field(default_factorylist) created_at: float field(default_factorytime.time) def to_dict(self) - Dict[str, Any]: return { task_id: self.task_id, task_type: self.task_type, payload: self.payload, history: self.history, created_at: self.created_at, }这里的关键点是把history放进任务对象。因为多智能体协作不是简单的一次请求响应而是多个节点依次处理。每个 Agent 做了什么、结果是什么都记录到history方便排查问题。再定义结果。结果不只是一个success布尔值还要包含数据和下一个阶段所需的负载# agency/models.py dataclass class TaskResult: success: bool data: Dict[str, Any] field(default_factorydict) error: str next_payload: Dict[str, Any] field(default_factorydict)其中next_payload是调度器需要关注的重要字段。Agent 不需要知道下一个 Agent 是谁它只需要把希望传给下一阶段的数据放到next_payload里。具体路由到哪个 Agent由调度器根据 pipeline 配置决定。2.4 Agent 抽象基类每个 Agent 都需要一个统一接口。基础接口可以定义如下# agency/base.py from abc import ABC, abstractmethod from agency.models import Task, TaskResult class Agent(ABC): name: str def __init__(self, timeout: int 10): self.timeout timeout abstractmethod def process(self, task: Task) - TaskResult: 处理任务并返回结果 raise NotImplementedError这里不建议把 Agent 的process方法直接写成异步协程。原因是最小示例里要先把调用链跑通异步可以放在调度层再加。实际项目中如果你的业务有大量 IO 等待可以把process改为async def process调度器再用asyncio或任务队列管理。后续会单独说明。2.5 注册中心和调度器注册中心负责保存所有 Agent 实例提供按名称获取的能力# agency/registry.py from typing import Dict, Optional from agency.base import Agent class AgentRegistry: def __init__(self): self._agents: Dict[str, Agent] {} def register(self, agent: Agent) - None: self._agents[agent.name] agent def get(self, name: str) - Optional[Agent]: return self._agents.get(name) def all_names(self): return list(self._agents.keys())调度器负责按 pipeline 顺序执行。这里使用线程池控制并发并用future.result(timeout)限制单个 Agent 的最大执行时间# agency/dispatcher.py import logging from concurrent.futures import ThreadPoolExecutor from typing import List, Optional from agency.models import Task, TaskResult from agency.registry import AgentRegistry from agency.base import Agent logger logging.getLogger(__name__) class Dispatcher: def __init__(self, registry: AgentRegistry, pipeline: List[str], max_workers: int 4): self.registry registry self.pipeline pipeline self.max_workers max_workers def run(self, task: Task) - TaskResult: current_payload task.payload with ThreadPoolExecutor(max_workersself.max_workers) as executor: for agent_name in self.pipeline: agent self.registry.get(agent_name) if agent is None: return TaskResult( successFalse, errorfagent not found: {agent_name} ) sub_task Task( task_idtask.task_id, task_typetask.task_type, payloadcurrent_payload, historytask.history ) logger.info(dispatch task %s to agent %s, task.task_id, agent_name) try: future executor.submit(agent.process, sub_task) result future.result(agent.timeout) except TimeoutError: return TaskResult( successFalse, errorfagent {agent_name} timeout ) except Exception as exc: logger.exception(agent %s failed, agent_name) return TaskResult( successFalse, errorfagent {agent_name} error: {exc} ) if not result.success: return result task.history.append({ agent: agent_name, data: result.data, error: result.error }) current_payload result.next_payload return TaskResult( successTrue, data{history: task.history}, next_payloadcurrent_payload )这里有一个重要设计为什么要复制一个sub_task而不是直接改task.payload因为任务对象可能被多个调度器或重试逻辑引用复制可以避免不同分支之间相互污染。实际项目中如果任务对象在数据库中持久化还要考虑版本号或乐观锁。3. 用两个 Agent 跑通一次端到端协作任务3.1 实现清洗 Agent先实现一个清洗文本的 Agent。它只负责把输入中的多余空白去掉并返回下一阶段需要的文本# agents/clean_agent.py from agency.base import Agent from agency.models import Task, TaskResult class CleanAgent(Agent): name clean_agent def process(self, task: Task) - TaskResult: raw_text task.payload.get(raw_text, ) cleaned_text .join(raw_text.split()) return TaskResult( successTrue, data{cleaned_text: cleaned_text}, next_payload{cleaned_text: cleaned_text} )这里data和next_payload都带有cleaned_text是因为data用于记录历史next_payload用于传给下一阶段。生产环境中两者可以完全一致也可以不同但最好保证每个 Agent 输出的next_payload字段名稳定避免下游 Agent 频繁改动。3.2 实现汇总 Agent第二个 Agent 负责根据清洗后的文本生成一个摘要结果。这里的摘要很粗糙只是为了演示数据传递# agents/summary_agent.py from agency.base import Agent from agency.models import Task, TaskResult class SummaryAgent(Agent): name summary_agent def process(self, task: Task) - TaskResult: cleaned_text task.payload.get(cleaned_text, ) summary fsummary of {len(cleaned_text)} chars return TaskResult( successTrue, data{summary: summary}, next_payload{summary: summary} )真实项目中摘要 Agent 可能会调用大模型、规则库或搜索引擎。但无论内部逻辑多复杂对外输出的结构都应该保持一致这样调度器不需要关心 Agent 内部实现。3.3 用 YAML 描述编排规则把 pipeline 放在 YAML 中而不是写死在代码里# config/pipeline.yaml pipeline: - clean_agent - summary_agent agents: clean_agent: timeout: 10 summary_agent: timeout: 5读取配置时可以使用简单的方法。如果项目引入了 PyYAML可以用下面的方式加载pip install pyyamlimport yaml with open(config/pipeline.yaml, r, encodingutf-8) as f: config yaml.safe_load(f)如果你的环境没有 PyYAML也可以先用 JSON 写配置结构是一样的。3.4 编写主程序入口主程序要做的事情初始化注册中心、注册 Agent、读取 pipeline、创建任务、交给调度器执行。# main.py import logging import yaml from agency.models import Task from agency.registry import AgentRegistry from agency.dispatcher import Dispatcher from agents.clean_agent import CleanAgent from agents.summary_agent import SummaryAgent logging.basicConfig( levellogging.INFO, format%(asctime)s %(levelname)s %(message)s ) def load_pipeline(path: str): with open(path, r, encodingutf-8) as f: config yaml.safe_load(f) return config[pipeline], config.get(agents, {}) def main(): pipeline, agent_config load_pipeline(config/pipeline.yaml) registry AgentRegistry() registry.register(CleanAgent(timeoutagent_config[clean_agent].get(timeout, 10))) registry.register(SummaryAgent(timeoutagent_config[summary_agent].get(timeout, 5))) dispatcher Dispatcher( registryregistry, pipelinepipeline, max_workers4 ) task Task( task_typefeedback_clean, payload{raw_text: Hello world , this is a test. } ) result dispatcher.run(task) if result.success: print(SUCCESS) print(result.data) else: print(FAILED) print(result.error) if __name__ __main__: main()这里把 Agent 的超时时间从配置读取便于后续调整。但要注意如果配置里没有对应的 key.get()的默认值要兜住避免 Agent 初始化失败。3.5 运行和预期输出执行下面的命令python main.py预期输出大致如下INFO dispatcher: dispatch task 8f3b2a1c to agent clean_agent INFO dispatcher: dispatch task 8f3b2a1c to agent summary_agent SUCCESS {history: [ {agent: clean_agent, data: {cleaned_text: Hello world , this is a test.}, error: }, {agent: summary_agent, data: {summary: summary of 29 chars}, error: } ]}这个结果看起来简单但它证明了几个关键点任务被顺序分发到两个 Agent中间结果被正确传递历史记录完整。如果后续业务需要增加一个敏感词过滤 Agent只需要在pipeline中插入名称然后注册一个 Agent主程序不需要大幅改动。4. 关键机制深入消息路由、状态流转和失败恢复4.1 消息字段设计多智能体项目最容易乱的地方是每个 Agent 各自定义输入输出字段最后字段名对不上。建议在项目初期就定义一套约定task_id一次协作任务的唯一标识所有 Agent 共享。task_type任务类型用于路由到不同 pipeline。payload当前阶段的有效负载。next_payload当前 Agent 传递给下一阶段的数据。history历史中间结果只追加不覆盖。status当前状态例如pending、running、success、failed。字段约束可以在models.py中使用 Pydantic 或 dataclass 校验。示例中使用 dataclass已经能保证最基本的类型正确。如果使用 Pydantic还可以增加字段长度限制和必填校验适合生产环境。4.2 状态流转一个任务在agency-agents中的状态大致如下pending - running - success | v failed其中running内部还可以细分dispatch、processing、waiting_result。这样设计的好处是当任务卡住时可以快速判断是调度器没分发出去还是 Agent 执行时间过长还是等待队列堆积。在调度器中每个环节都应该有状态记录。如果只用线程池submit一旦线程池线程不够任务会排队此时状态需要从pending转为queued再转为running。最小示例没有单独标记但生产项目必须区分。状态含义常见卡点pending任务已创建等待调度调度器未启动queued线程池或队列中等待并发数耗尽runningAgent 正在处理Agent 内部 IO 阻塞success所有 pipeline 节点完成无failed某个 Agent 失败或超时需要回滚或重试4.3 失败重试策略Agent 失败是常态不能把失败当作异常直接抛出。常见策略有可重试的失败网络抖动、依赖服务超时设置重试次数。不可重试的失败参数错误、数据格式不对直接返回失败。部分失败某个 Agent 处理后返回失败选择终止整条链路或跳过。在调度器中最简单的方式是给每个 Agent 配置retriesagents: clean_agent: timeout: 10 retries: 2 summary_agent: timeout: 5 retries: 1调度器在处理超时或异常时根据retries决定是否重新执行# 伪代码示意 for retry_count in range(retries 1): try: result future.result(agent.timeout) if result.success: break except TimeoutError: logger.warning(agent %s timeout, retry %s, agent_name, retry_count) if retry_count retries: return TaskResult(successFalse, errortimeout)注意重试前要考虑幂等性。如果 Agent 内部有写库或发送消息的动作重试可能会导致重复写入。所以每个 Agent 在设计时就要遵循“处理任务一次结果一致”的幂等原则或者至少记录任务 ID在收到重复任务时直接返回已有结果。4.4 并发调优参数调度并发数、超时时间、队列大小是影响整个系统吞吐的关键参数。max_workers线程池最大线程数。太小会排队太大会增加资源竞争。建议从 CPU 核数开始调。timeout每个 Agent 的超时时间。不能统一设置因为清洗任务和模型推理的耗时差异很大。queue_size如果引入任务队列队列大小要能缓冲业务高峰避免消息直接丢弃。如果 Agent 主要做计算型任务max_workers可以设置为CPU 核数 1。如果 Agent 主要做 IO 操作比如调用 HTTP API可以适当调大但要注意下游服务的限流。对于外部 API 调用超时时间要设置得保守一些否则线程会被长时间占住。5. 常见问题排查任务卡住、重复执行和结果不一致5.1 任务卡住没有日志现象调度器已经打印了dispatch task xxx to agent yyy然后就没有下一步日志。排查步骤检查 Agent 内部是否存在sleep或线程阻塞。检查是否调用了外部接口且外部接口没有超时设置。检查agent.timeout是否配置过小导致线程池中的线程长期被占住。检查线程池是否被占满其他任务在排队。解决方式给每个外部调用设置显式超时避免无限等下去。对线程池队列长度加监控。日志里除了指示“开始调度”还要在 Agent 执行前后打印耗时start time.time() result agent.process(sub_task) logger.info(agent %s finished, cost%.2fs, agent_name, time.time() - start)这样可以快速定位是哪个 Agent 耗时高。5.2 任务被重复执行现象同一task_id出现在多条执行日志中或者下游系统收到重复数据。原因Agent 内部有重试逻辑而调度器又做了一层重试。消费方没有做去重。任务超时但实际上 Agent 还在执行调度器把这个任务重新提交到另一个线程。排查方式在流程入口根据task_id查询是否已存在。在 Agent 执行结果中加入执行批次号。检查项目中有几层重试逻辑避免叠加。推荐方案是在任务入库时以task_id做唯一索引并且 Agent 保存结果时使用乐观锁。若检测到任务状态已经是success直接忽略重复请求。5.3 Agent 之间的结果对不上现象A 输出字段是cleaned_textB 却读取clean_text导致变量为空。原因字段命名在不同 Agent 间没有统一约束或者某个 Agent 升级后改变了输出结构。排查方式查看history中每个 Agent 的data。检查 pipeline 配置是否与代码一致。在调度器中增加字段校验对每个 Agent 的next_payload做 schema 校验。推荐在models.py中为每个阶段定义独立的 payload 类型。这样 A Agent 的输出会先被转换为下一阶段期望的类型字段不匹配时系统能提前报错。5.4 并发死锁现象任务全部卡住线程池没有线程处理。可能是一个 Agent 内部又调用了同一个调度器等待某个子任务完成而子任务需要等待该 Agent 释放线程形成互相等待。排查方式打印线程堆栈。检查是否有 Agent 内部调用dispatcher.run。使用独立线程池或异步任务避免同步嵌套。解决方式不要让单个 Agent 直接调度整个 pipeline。如果需要子任务应当把子任务发送到独立队列或者使用异步消息系统。最小示例中ThreadPoolExecutor内部如果再次调用Dispatcher.run很容易把线程池占满。5.5 排查清单问题现象可能原因检查方式处理建议任务卡住Agent 未设置超时或外部调用阻塞看线程堆栈、Agent 耗时日志增加超时和熔断任务重复执行多层重试叠加检查调度器和 Agent 重试配置以 task_id 做幂等结果字段不匹配缺少 schema 校验查看 history 数据定义阶段化 payload并发死锁Agent 内部嵌套调度打印线程堆栈拆分子任务队列Agent 未注册配置不一致打印 registry.all_names()启动时自动校验配置配置不生效修改了错误环境检查加载路径和环境变量启动时打印配置摘要6. 从学习环境到生产环境注册中心、可观测性和发布回滚6.1 配置外置化最小示例中pipeline 在本地 YAML 文件里。测试环境可以直接读取本地文件但生产环境通常要把配置放到配置中心或环境变量中。这样调整 Agent 列表、修改超时和重试次数时不需要重新发版。推荐做法将业务配置放在 YAML 或 JSON 中启动时加载。将密钥、数据库地址等放到环境变量或配置中心。发布前校验配置避免 pipeline 指向不存在的 Agent。如果配置发生变化要有版本记录和回滚机制。6.2 日志与可观测性多智能体项目最需要的是链路追踪。每个任务要有统一的task_id并且日志中始终携带这个 ID。否则同一任务被多个 Agent 处理时日志会散落在不同地方很难串联。可以在日志格式中加上task_id2025-05-10 15:33:22 INFO [task_id8f3b2a1c] dispatch to clean_agent 2025-05-10 15:33:22 INFO [task_id8f3b2a1c] clean_agent finished 2025-05-10 15:33:23 INFO [task_id8f3b2a1c] dispatch to summary_agent生产环境还可以把每个 Agent 的执行耗时、重试次数和失败原因输出到指标系统。对于外部 API 调用要记录调用参数脱敏后的摘要、状态码和耗时。6.3 隔离与权限如果多个业务线共用一套agency-agents建议按业务线隔离 Agent 注册表而不是全部放在一个注册中心里。隔离可以是进程级也可以是命名空间级。权限方面要注意Agent 之间不能随意注册和替换实现。必须限制外部系统对 Agent 注册表 API 的访问。每个 Agent 执行的最小权限原则不能让一个文本清洗 Agent 拥有数据库删除权限。6.4 发布与回滚多智能体系统上线前除了常规测试还要做 pipeline 变更演练。例如新 Agent 上线先灰度 10% 的流量。修改调度器并发参数先压测再全量。如果某个 Agent 出现问题可以快速从 pipeline 中摘除而不是紧急改代码。发布检查清单所有 Agent 已注册配置中的 pipeline 都能找到对应实现。每个 Agent 的超时、重试参数合理。关键 Agent 有幂等保护。日志中能看到完整链路和中间结果。有回滚方案例如保留上一份配置快照。监控大盘能看到任务成功率和平均耗时。对下游系统的调用有熔断和降级机制。6.5 什么时候不应该自己造轮子agency-agents最小示例可以帮助理解原理但生产级多智能体系统往往需要以下能力持久化任务状态。分布式调度。可靠消息队列。人工审核和干预。可视化编排。如果项目复杂度已经超过“几个 Agent 顺序调用”优先评估现有开源方案。使用成熟框架时仍然要理解它的任务模型、消息结构和状态流转这比背 API 重要得多。自己实现时最值得投入的也不是 Agent 内部逻辑而是把任务生命周期、幂等和可观测性设计好。多智能体协作系统的工程难点不在于单个 Agent 写得多么复杂而在于如何把多个 Agent 组织成一条可维护、可追踪、可恢复的流水线。agency-agents这个名字背后代表的就是这样一套组织方式agency 负责调度和管理agents 负责专注执行。先把任务模型、注册中心、调度器三个核心组件搭好再逐步加入配置、重试和监控多智能体项目才能从玩具变为可交付的生产系统。如果你正在设计自己的多智能体项目建议从最小串联链路开始先跑通两个 Agent再扩展更多节点并从一开始就记录每一步的中间结果。这样做最大的价值是当系统出问题时你能从日志和历史记录中准确知道哪一步出了问题。