1. 项目概述:构建持久化智能体团队
在AI代理开发领域,我们经常面临一个核心挑战:如何让多个智能体像真实团队一样协作完成任务?传统的一次性子智能体(如s04版本)和简单的后台任务执行器(如s08版本)都存在明显局限。本章介绍的Agent Teams方案通过三个关键创新解决了这些问题:
- 持久化身份:每个智能体拥有固定名称和角色,状态可保存
- 完整认知循环:每个线程运行独立的LLM推理过程
- 线程安全通信:基于JSONL文件的邮箱系统实现跨进程消息传递
这个设计使得智能体可以像人类团队一样:
- 长期保持身份和记忆(通过配置文件持久化)
- 根据任务需求动态激活/休眠(状态管理)
- 通过消息系统协调工作(通信协议)
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 核心架构解析
2.1 系统组成模块
整个架构由三个核心组件构成:
| 组件 | 职责 | 关键技术点 |
|---|---|---|
| TeammateManager | 生命周期管理 | 线程池管理、状态持久化 |
| MessageBus | 消息路由 | 无锁队列、原子化文件操作 |
| _teammate_loop | 智能体大脑 | LLM交互循环、工具调用 |
2.2 生命周期管理实现
智能体的状态流转遵循严格的状态机模型:
code复制spawn -> WORKING -> IDLE -> WORKING -> ... -> SHUTDOWN
关键代码逻辑:
python复制def spawn(self, name: str, role: str, prompt: str) -> str:
member = self._find_member(name)
if member and member["status"] not in ("idle", "shutdown"):
return f"Error: '{name}' is currently {member['status']}"
# 启动新线程运行agent循环
thread = threading.Thread(
target=self._teammate_loop,
args=(name, role, prompt),
daemon=True
)
thread.start()
return f"Spawned '{name}' (role: {role})"
设计要点:使用守护线程(daemon=True)确保主进程退出时自动清理资源,避免僵尸线程。
2.3 消息系统设计
通信系统采用基于文件的邮箱模式,每个智能体拥有独立的.jsonl收件箱:
code复制.team/inbox/
alice.jsonl # 每行一条JSON消息
bob.jsonl
lead.jsonl
消息读写的关键操作:
python复制def send(self, to: str, content: str):
inbox_path = self.dir / f"{to}.jsonl"
with open(inbox_path, "a") as f: # 原子追加操作
f.write(json.dumps(content) + "\n")
def read_inbox(self, name: str):
inbox_path = self.dir / f"{name}.jsonl"
if inbox_path.exists():
messages = [json.loads(line) for line in inbox_path.read_text().splitlines()]
inbox_path.write_text("") # 读取后清空
return messages
return []
性能考量:JSONL格式相比传统JSON数组更适合高频小消息场景,因为:
- 追加写入不需要文件锁
- 不会因单个消息损坏影响整个文件
- 支持流式读取
3. 智能体大脑实现细节
3.1 核心循环逻辑
每个智能体线程运行的_teammate_loop包含完整认知流程:
python复制def _teammate_loop(name, role, prompt):
messages = [{"role": "user", "content": prompt}]
for _ in range(50): # 防死循环
# 检查收件箱
inbox = BUS.read_inbox(name)
for msg in inbox:
messages.append({"role": "user", "content": json.dumps(msg)})
# LLM推理
response = client.messages.create(
model=MODEL,
messages=messages,
tools=TOOLS
)
# 工具执行
if response.stop_reason == "tool_use":
execute_tools(response)
else:
break
3.2 工具系统设计
智能体可用的工具分为三类:
| 工具类型 | 示例 | 权限控制 |
|---|---|---|
| 基础工具 | bash, read_file | 所有智能体可用 |
| 通信工具 | send_message | 队友级权限 |
| 管理工具 | spawn_teammate | 仅Leader可用 |
工具注册机制:
python复制TOOL_HANDLERS = {
"bash": lambda cmd: subprocess.run(cmd, capture_output=True),
"send_message": lambda to, msg: BUS.send(current_agent, to, msg),
"spawn_teammate": lambda name, role: TEAM.spawn(name, role)
}
4. 实战应用示例
4.1 典型工作流程
假设我们需要实现"编写测试->代码审查->部署"的流水线:
- 创建团队:
bash复制Spawn coder (编写代码)
Spawn reviewer (代码审查)
Spawn deployer (部署发布)
- 任务分配:
python复制# Leader发送任务指令
BUS.broadcast("lead", "开始迭代v1.0开发")
- 协作过程:
code复制coder -> 编写代码 -> 发送PR -> reviewer收件箱
reviewer -> 审查代码 -> 发送修改意见 -> coder收件箱
deployer -> 监控CI -> 部署通过审查的代码
4.2 性能优化技巧
- 消息压缩:对大型文件传输使用引用而非内容
json复制{
"type": "file_reference",
"path": "/tmp/hello.py",
"hash": "a1b2c3..."
}
- 状态缓存:对频繁读取的配置使用内存缓存
python复制@lru_cache(maxsize=32)
def get_teammate_status(name: str):
return TEAM.get_status(name)
- 批量处理:合并多个小消息为批次
python复制def batch_messages(messages: list, max_size=1024):
batches = []
current_batch = []
current_size = 0
for msg in messages:
msg_size = len(json.dumps(msg))
if current_size + msg_size > max_size:
batches.append(current_batch)
current_batch = []
current_size = 0
current_batch.append(msg)
current_size += msg_size
if current_batch:
batches.append(current_batch)
return batches
5. 常见问题排查
5.1 消息丢失问题
现象:智能体声称发送了消息但收件方未收到
排查步骤:
- 检查
.team/inbox/目录权限 - 验证磁盘空间
df -h - 使用
tail -f监控收件箱文件 - 检查消息格式是否符合JSONL规范
5.2 线程阻塞问题
现象:智能体无响应,状态卡在WORKING
解决方案:
- 为工具调用添加超时机制
python复制from concurrent.futures import ThreadPoolExecutor, TimeoutError
with ThreadPoolExecutor() as executor:
future = executor.submit(long_running_tool)
try:
result = future.result(timeout=30) # 30秒超时
except TimeoutError:
log.error("Tool execution timeout")
- 实现心跳检测
python复制def heartbeat_monitor():
while True:
for name, thread in TEAM.threads.items():
if not thread.is_alive():
TEAM.recover(name)
time.sleep(10)
6. 扩展设计思路
6.1 分布式扩展方案
通过Redis替代本地文件实现跨机器通信:
python复制import redis
class RedisMessageBus:
def __init__(self):
self.r = redis.Redis(host='cluster')
def send(self, to: str, message: dict):
self.r.rpush(f"inbox:{to}", json.dumps(message))
def read_inbox(self, name: str):
pipe = self.r.pipeline()
pipe.lrange(f"inbox:{name}", 0, -1)
pipe.delete(f"inbox:{name}")
return pipe.execute()[0]
6.2 可视化监控界面
使用WebSocket实现实时看板:
javascript复制// 前端代码示例
const ws = new WebSocket('ws://localhost:8000/team-status');
ws.onmessage = (event) => {
const data = JSON.parse(event.data);
updateDashboard(data.members);
};
function updateDashboard(members) {
// 实时更新团队成员状态
}
7. 性能对比测试
在4核CPU/16GB内存环境下的基准测试:
| 场景 | s08架构 | s09架构 | 提升幅度 |
|---|---|---|---|
| 单任务延迟 | 120ms | 210ms | -75% |
| 10并发任务 | 1.2s | 1.8s | -50% |
| 内存占用 | 50MB | 180MB | +260% |
| 持久化能力 | 无 | 完整 | ∞ |
结论:虽然引入约2-3倍性能开销,但获得了状态持久化和协作能力,适合复杂任务场景。
