多Agent协作的Python实现:从零构建Swarm-forge协调器
2026/8/30 22:51:24 网站建设 项目流程

当多个 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 协调器的核心职责可以收敛为四点

一个简单协调器不需要一开始就做成完整框架。按照最少可用原则,可以把职责收敛成四件事:

  1. 注册:让 agent 提供自己的名称、描述和执行函数。
  2. 分发:把一个任务发给匹配的 agent,或者按配置指定 agent。
  3. 执行:支持串行、并行以及带超时和重试的执行。
  4. 聚合:收集每个 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示例过程不强制需要,可离线运行

学习环境下,用标准库的queuethreadingdataclasses就能把协调器跑通。如果一开始就引入 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_factory=lambda: 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_id=task.task_id, agent_name=task.agent_name, status="error", output=None, started_at=now, finished_at=now, error=error, )

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(f"agent 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(f"agent 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_id=uuid.uuid4().hex, agent_name=agent_name, payload=payload, timeout=timeout or self.task_timeout, retries=retries 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( target=self._run_loop, daemon=True, name="swarm-worker", ) t.start() self._workers.append(t) def shutdown(self) -> None: self._stop_event.set() for t in self._workers: t.join(timeout=1.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(timeout=0.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, error=str(exc))) return if self._stop_event.wait(timeout=0.5): self._save_result( AgentResult.from_error(task, error="stopped during retry") ) return def _execute_once(self, task: Task) -> AgentResult: agent = self.registry.get(task.agent_name) if not isinstance(agent, Agent): raise TypeError(f"registered 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_id=task.task_id, agent_name=task.agent_name, status="success", output=output, started_at=started_at, finished_at=finished_at, error=None, ) 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_workers=2, task_timeout=10, max_retries=1) 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中包含预期字段。如果某个statuserror,需要查看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 支持

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询