十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

多Agent协作的Python实现:从零构建Swarm-forge协调器

多Agent协作的Python实现:从零构建Swarm-forge协调器 当多个 AI agent 需要协作完成同一件任务时最直接的做法是让每个 agent 单独处理一个子任务再由一个调度者统一收集结果。Swarm-forge 就是围绕这个需求设计的简单工具它不负责训练模型也不负责具体业务逻辑只负责把多个 agent 的注册、调用、并发调度和结果汇总变成一套可复用流程。这篇文章会带你在 Python 中从零实现一个 Swarm-forge 风格的最小协调器。代码量不大但足以讲清多 agent 编排的核心链路。读完你能够理解 agent 注册表、任务队列、结果聚合在协调器中分别承担什么职责并能把这个最小实现作为起点扩展成自己项目里的多 agent 调度基础。1. 理解多 Agent 协调要解决什么问题1.1 为什么单 Agent 不够用很多实际任务不是一次提示词就能完成的。比如“写一篇技术文章并检查可读性”通常需要先生成草稿再让另一个角色从逻辑、语法、信息密度角度评审。如果只用一个 agent 串行完成整个流程是写死的如果要调整评审角色、更换模型、增加并行检查代码很容易变成一堆 if/else 分支。单 agent 的瓶颈可以归纳为三点上下文限制把全部资料塞进一个 agent 的上下文容易超过模型窗口也会让 prompt 越来越难维护。职责耦合生成、总结、质检、翻译往往需要不同的 prompt 和不同的模型参数硬写在一个函数中后续改动成本很高。并发困难一个 agent 内部写串行逻辑容易但要同时处理多个资源、多个角色需要额外的调度能力。多 agent 协作的核心思路是拆分每个 agent 只负责一个小而明确的任务协调器负责把它们组织成完整执行流。1.2 Swarm-forge 在协调链路中的位置Swarm-forge 名字里有两个关键词Swarm 表示多个 agent 组成的群体forge 表示把这些分散部件锻造成一条可运行的链路。它在整体架构中位于上层业务和底层模型 API 之间。实际项目里可以这样分层业务层用户请求、文件上传、最终结果展示。协调层Swarm-forge 负责任务拆分、agent 注册、调度、重试、结果聚合。Agent 执行层每个 agent 内部完成 prompt 组装、调用模型、解析响应。基础设施层模型 API、数据库、消息队列、日志存储。Swarm-forge 的价值在于让上层业务只依赖协调器接口而不需要知道每个 agent 内部是怎么实现的。也可以反过来理解如果项目里只有一次模型调用不需要引入多 agent 协调器当业务开始出现按角色拆分、按子任务并行、按结果串联的需求时才值得把协调逻辑单独抽出来。1.3 协调器的核心职责可以收敛为四点一个简单协调器不需要一开始就做成完整框架。按照最少可用原则可以把职责收敛成四件事注册让 agent 提供自己的名称、描述和执行函数。分发把一个任务发给匹配的 agent或者按配置指定 agent。执行支持串行、并行以及带超时和重试的执行。聚合收集每个 agent 的执行结果整理为可读输出。这四条是后面所有代码实现的主线。后面章节里的类、函数和参数全部围绕这四个职责展开。2. 准备环境并设计一个极简协调器的模块边界2.1 环境要求实现 Swarm-forge 最小版本只需要 Python 3.9 及以上版本不需要第三方依赖。示例代码使用了dict[str, Any]这样自带泛型支持的注解所以 Python 版本不能太低。推荐在虚拟环境里测试避免污染系统环境。python -m venv .venv source .venv/bin/activate # Windows 使用 .venv\Scripts\activate环境要求如下表项目要求Python3.9 或更高第三方依赖无标准库即可操作系统Windows / Linux / macOS 均可模型 API示例过程不强制需要可离线运行学习环境下用标准库的queue、threading、dataclasses就能把协调器跑通。如果一开始就引入 Celery、Redis、Kafka 等组件反而会掩盖协调器本身的代码逻辑。2.2 三个核心模块的职责根据前面收敛的四个职责设计三个核心模块AgentRegistry保存所有已注册的 agent。最简实现可以用字典key 是 agent 名称value 是 Agent 对象。TaskQueue保存等待执行的任务。单机版直接用queue.Queue分布式版本可以换成 Redis Stream、RabbitMQ 或 Kafka。ResultAggregator收集执行结果。为了支持按任务 ID 回溯可以用字典保存。它们的关系是外部提交任务协调器把任务放入队列工作线程从队列取出任务根据 agent 名称从注册表获取执行器执行后把结果写入聚合器。这里不需要一开始就引入 DAG 调度。DAG 适合处理复杂的任务依赖但会增加很多概念。作为学习版本先用队列模型把主流程讲清楚后续再扩展依赖关系。2.3 项目目录结构为了让代码边界清晰把不同职责拆到不同文件swarm_forge/ ├── __init__.py # 导出核心类 ├── agent.py # Agent 抽象基类 ├── registry.py # AgentRegistry 注册表 ├── task.py # Task 和 AgentResult 数据模型 ├── forge.py # SwarmForge 协调器 └── config.py # 配置加载 examples/ ├── content_agents.py └── run_example.py单一文件也能实现同样功能但拆文件以后你要增加分布式队列、自定义 Agent、配置中心时不需要改动已有接口。Agent抽象类是最重要的接口约定协调器只依赖run(payload)方法。3. 从零实现 Swarm-forge 核心代码3.1 定义 Task 与 Agent 抽象先定义任务和结果的数据模型。任务需要唯一 ID这样才能在并发执行后准确聚合结果。# task.py from dataclasses import dataclass, field from datetime import datetime, timezone from typing import Any, Optional dataclass class Task: task_id: str agent_name: str payload: dict[str, Any] timeout: float 30.0 retries: int 0 created_at: str field( default_factorylambda: datetime.now(timezone.utc).isoformat() ) dataclass class AgentResult: task_id: str agent_name: str status: str output: Optional[dict[str, Any]] started_at: str finished_at: str error: Optional[str] None classmethod def from_error(cls, task: Task, error: str) - AgentResult: now datetime.now(timezone.utc).isoformat() return cls( task_idtask.task_id, agent_nametask.agent_name, statuserror, outputNone, started_atnow, finished_atnow, errorerror, )Task里的payload是任意结构化字典AgentResult里的output也是字典。这样设计的好处是方便 JSON 序列化后续如果要把任务和结果写入数据库不需要做复杂转换。Agent 抽象接口# agent.py from abc import ABC, abstractmethod from typing import Any class Agent(ABC): agent_name: str description: str abstractmethod def run(self, payload: dict[str, Any]) - dict[str, Any]: raise NotImplementedError协调器不需要知道 agent 内部是调用大模型、执行函数还是查数据库只需要统一入口。真实项目里还可以增加async run版本但学习版本先用同步实现避免并发模型干扰主逻辑。3.2 实现 AgentRegistry注册表解决的是“根据名字找到 agent”的问题。# registry.py from typing import Dict from swarm_forge.agent import Agent class AgentRegistry: def __init__(self) - None: self._agents: Dict[str, Agent] {} def register(self, agent: Agent) - None: if agent.agent_name in self._agents: raise ValueError(fagent already exists: {agent.agent_name}) self._agents[agent.agent_name] agent def get(self, name: str) - Agent: try: return self._agents[name] except KeyError as exc: raise KeyError(fagent not found: {name}) from exc def list_agents(self) - list[str]: return sorted(self._agents.keys())这里做“名称唯一”校验是为了防止多个同名 agent 被静默覆盖。生产环境里同名覆盖很容易引发线上事故比如注册了两个 writer但后一个覆盖前一个调用结果完全不可控。3.3 实现 SwarmForge 协调器协调器是核心。它负责接收任务、启动 worker 线程、执行 agent、保存结果。# forge.py import queue import threading import uuid from datetime import datetime, timezone from typing import Any, Optional from swarm_forge.agent import Agent from swarm_forge.registry import AgentRegistry from swarm_forge.task import AgentResult, Task class SwarmForge: def __init__( self, registry: AgentRegistry, max_workers: int 4, task_timeout: float 30.0, max_retries: int 0, ) - None: self.registry registry self.max_workers max_workers self.task_timeout task_timeout self.max_retries max_retries self._task_queue: queue.Queue[Task] queue.Queue() self._results: dict[str, AgentResult] {} self._lock threading.Lock() self._stop_event threading.Event() self._workers: list[threading.Thread] [] self._started False def submit( self, agent_name: str, payload: dict[str, Any], timeout: Optional[float] None, retries: Optional[int] None, ) - str: task Task( task_iduuid.uuid4().hex, agent_nameagent_name, payloadpayload, timeouttimeout or self.task_timeout, retriesretries if retries is not None else self.max_retries, ) self._task_queue.put(task) return task.task_id def start(self) - None: if self._started: return self._started True for _ in range(self.max_workers): t threading.Thread( targetself._run_loop, daemonTrue, nameswarm-worker, ) t.start() self._workers.append(t) def shutdown(self) - None: self._stop_event.set() for t in self._workers: t.join(timeout1.0) self._started False self._workers.clear() def _run_loop(self) - None: while not self._stop_event.is_set(): try: task self._task_queue.get(timeout0.5) except queue.Empty: continue try: self._execute_with_retry(task) finally: self._task_queue.task_done() def _execute_with_retry(self, task: Task) - None: attempt 0 while True: attempt 1 try: result self._execute_once(task) self._save_result(result) return except Exception as exc: if attempt task.retries: self._save_result(AgentResult.from_error(task, errorstr(exc))) return if self._stop_event.wait(timeout0.5): self._save_result( AgentResult.from_error(task, errorstopped during retry) ) return def _execute_once(self, task: Task) - AgentResult: agent self.registry.get(task.agent_name) if not isinstance(agent, Agent): raise TypeError(fregistered object is not Agent: {task.agent_name}) started_at datetime.now(timezone.utc).isoformat() output agent.run(task.payload) finished_at datetime.now(timezone.utc).isoformat() return AgentResult( task_idtask.task_id, agent_nametask.agent_name, statussuccess, outputoutput, started_atstarted_at, finished_atfinished_at, errorNone, ) def _save_result(self, result: AgentResult) - None: with self._lock: self._results[result.task_id] result def wait(self, timeout: Optional[float] None) - dict[str, AgentResult]: self._task_queue.join() with self._lock: return dict(self._results)几个关键点submit()负责生成任务 ID并放入队列。任务 ID 是后续查询结果的依据。start()启动固定数量的 worker 线程。线程数量就是并发度。_run_loop()不断从队列取任务并在finally中调用task_done()这样wait()才能通过队列的join()判断全部任务完成。_save_result()使用锁保护共享字典避免多个 worker 线程同时写结果导致数据丢失。task_timeout在这个最小版本中只作为配置字段保存并没有真正中断已经卡死的 agent。真正严格的超时要依赖Future.result(timeout...)或子进程隔离需要结合你实际采用的 agent 执行方式实现。3.4 配置入口与导出简单配置可以从环境变量读取方便在命令行临时调整。# config.py import os def load_config() - dict: return { max_workers: int(os.getenv(SWARM_MAX_WORKERS, 4)), task_timeout: float(os.getenv(SWARM_TASK_TIMEOUT, 30)), max_retries: int(os.getenv(SWARM_MAX_RETRIES, 0)), }在包的__init__.py中导出核心类调用方 import 起来更简洁# __init__.py from swarm_forge.agent import Agent from swarm_forge.registry import AgentRegistry from swarm_forge.forge import SwarmForge from swarm_forge.task import AgentResult, Task __all__ [Agent, AgentRegistry, AgentResult, SwarmForge, Task]配置代码虽然短但体现了一个原则不要把所有参数硬编码在业务文件里。学习环境可以用环境变量生产环境建议改成 YAML 或配置中心但对外接口要保持一致。4. 跑通一个双 Agent 协作的最小示例4.1 示例需求先生成再评审示例目标一个 writer agent 生成产品描述一个 reviewer agent 对描述进行评审。先用 writer 生成再把 writer 的输出作为 reviewer 的输入形成一次串行依赖。这个例子能验证注册、分发、执行、聚合全链路。为了在没有模型 API Key 的环境下也能运行示例直接用模拟结果代替真实模型调用。真实项目里只需要在run()方法中把模拟逻辑换成 prompt 组装和 API 调用。4.2 编写两个 Agent 类# examples/content_agents.py from swarm_forge.agent import Agent class WriterAgent(Agent): agent_name writer description 生成产品描述草稿 def run(self, payload: dict) - dict: product payload.get(product, default product) # 实际项目这里会组装 prompt 并调用模型 return { draft: f{product} 是一款面向日常场景的工具设计简洁使用成本低。 } class ReviewerAgent(Agent): agent_name reviewer description 检查文本长度并给出评审意见 def run(self, payload: dict) - dict: draft payload.get(draft, ) length len(draft) if length 20: opinion 内容太短需要补充细节。 else: opinion 内容长度合适建议补充使用场景。 return {length: length, opinion: opinion}评审 agent 的输入来自 writer 的输出所以payload中必须有draft字段。协调器本身不感知任务依赖依赖关系由调用方通过submit()顺序控制。4.3 运行脚本与预期输出# examples/run_example.py from swarm_forge import SwarmForge from swarm_forge.registry import AgentRegistry from examples.content_agents import ReviewerAgent, WriterAgent def main(): registry AgentRegistry() registry.register(WriterAgent()) registry.register(ReviewerAgent()) forge SwarmForge(registry, max_workers2, task_timeout10, max_retries1) forge.start() writer_task_id forge.submit(writer, {product: 便携蓝牙键盘}) results forge.wait() writer_result results[writer_task_id] reviewer_task_id forge.submit(reviewer, writer_result.output) results forge.wait() for task_id, result in results.items(): print(task_id, result.agent_name, result.status, result.output) forge.shutdown() if __name__ __main__: main()运行python examples/run_example.py预期输出类似1f3a9c2b0e0d4e5e8f6a2d3c4b5e6f7a writer success {draft: 便携蓝牙键盘 是一款面向日常场景的工具设计简洁使用成本低。} 8f7b6a5c4d3e2f1a0b9c8d7e6f5a4b3c reviewer success {length: 34, opinion: 内容长度合适建议补充使用场景。}验证标准所有status都是successoutput中包含预期字段。如果某个status是error需要查看result.error定位原因。5. 关键设计细节与参数说明5.1 并发度、超时和重试的作用边界max_workers是并发线程数。调大可以提高吞吐但也要看模型 API 的限流和内存占用。如果每个 agent 内部是 CPU 密集型处理线程数不是越大越好如果是 I/O 密集型模型调用适当地调大并发可以缩短整体耗时。task_timeout用来给单个任务设定预期上限。当前示例代码只保存了这个字段用于展示参数传递路径。真正要强制执行超时可以在_execute_once中使用concurrent.futures.Future.result(timeout...)但 agent 的run()必须能响应线程中断否则超时后线程仍会继续占用资源。更严格的隔离方案是把 agent 放进子进程执行。max_retries对暂时性故障有效比如网络抖动、API 限流。对于业务逻辑错误重试没有意义还可能重复提交。生产环境应该根据异常类型决定是否重试而不是对所有异常统一重试。5.2 任务依赖与结果传递方式当前示例展示了顺序依赖先执行 writer再把 writer 结果作为 reviewer 输入。实际场景中常见依赖模式有三种| 依赖模式 | 场景 | Swarm-forge 支持
返回列表