Appearance
Agent 服务拆分:将 LangGraph Agent 独立为无状态 API + Agent Worker 架构
背景与目标
要解决的问题
当前 backend/agent/ 中的三大核心 Agent(ChatAgent、ImageAgent、VideoAgent)与后端 FastAPI 服务运行在同一进程中。这导致:
- API 服务有状态:Agent 通过 LangGraph checkpointer 和上下文管理持有会话状态,API 服务无法真正无状态水平扩展。
- 资源竞争:LLM 推理的长时间连接和内存占用影响 API 服务的响应延迟和稳定性。
- 独立扩展受限:无法单独为 Agent 推理层增加计算资源(GPU、内存),必须整体扩容。
- 故障隔离弱:Agent 内部异常(如 LLM 超时、工具调用失败)可能影响 API 服务的其他请求。
可感知结果
- 后端 API 服务变为无状态服务,可自由水平扩展。
- Agent 推理作为独立 Worker 进程运行,可独立部署、扩容、回滚。
- 前端用户无感知:SSE 流式体验保持一致。
- 新架构复用现有 RocketMQ + Redis Stream 基础设施,与 RAG/Review Worker 模式统一。
现状分析
当前架构
Agent 层的外部依赖
| 依赖 | 来源 | 使用位置 | 拆分后归属 |
|---|---|---|---|
services.checkpointer | services/checkpointer.py | 三个 Agent 的 *_stream() 方法 | Agent Worker 自带(MySQL 连接) |
core.database | core/database.py | Image/Video Agent 的 _build_runtime_config() | Agent Worker 自带(读取推理配置) |
services.agent_config_service | services/agent_config_service.py | Image/Video Agent 的 _build_runtime_config() | Agent Worker 自带 |
services.web_search_service | services/web_search_service.py | agent/shared/web_search_tool.py | Agent Worker 自带 |
services.platform_skill_service | services/platform_skill_service.py | agent/shared/skill_runtime_resolver.py | 保留在 API 服务(由 Service 层解析后序列化传递) |
services.review.skill.service | services/review/skill/service.py | agent/shared/skill_runtime_resolver.py | 保留在 API 服务(同上) |
core.user_info.UserInfo | core/user_info.py | 三个 Agent 构建 LangSmith config | Agent Worker 自带(纯数据类) |
core.process_pool | core/process_pool.py | Video Agent 的 prompt_optimizer | Agent Worker 自带 |
| Provider 初始化 | core/providers.py | Image/Video tools 中调用 generate_image 等 | Agent Worker 自带 |
可复用的现有能力
- RocketMQ 集成:
rocketmq_integration/已封装 Producer/Consumer,支持自动重连和延迟初始化。 - Redis Stream 续接:
services/core/resume/已实现完整的 Stream 读写、并发控制、租约心跳机制。 - Worker 基础设施:
worker/bootstrap.py提供日志、DB 连接池、Redis、信号处理等共享启动逻辑。 - Docker 模式:
Dockerfile.rag和Dockerfile.review-agent已有独立 Worker 容器的构建模式。
限制
- Agent 间通过共享
thread_id和 MySQL checkpointer 协作,拆分后仍需共享同一数据库。 web_search_tool调用的搜索引擎 API 配置在 API 侧,Agent Worker 需要独立持有此配置。- SkillRuntimeResolver 依赖 DB 查询且依赖
services/层的服务类,不适合整体迁移。
第一性原理分析
核心业务对象
| 对象 | 职责 | 不变量 | 状态归属 |
|---|---|---|---|
| AgentRequest | 封装一次 Agent 调用的全部输入 | 必须包含 prompt、model、provider 凭证 | API 侧创建 → 序列化传递 |
| AgentResponse | 封装 Agent 产出的流式事件序列 | 事件按序递增,终态事件唯一 | Agent Worker 创建 → Redis Stream |
| AgentWorker | 消费 RocketMQ 消息、执行 Agent、产出事件 | 每条消息恰好处理一次(at-least-once + 幂等) | Agent Worker 进程内 |
| StreamBridge | API 侧桥接 Redis Stream → SSE | 客户端断线后可续接 | API 侧(已存在) |
| CheckpointerSession | LangGraph checkpoint 的数据库连接会话 | 并发连接数受限(当前上限 15) | Agent Worker 进程内 |
协作关系
变化点
| 变化维度 | 当前 | 拆分后 |
|---|---|---|
| Agent 类型 | Chat / Image / Video 三种 | 同,通过 RocketMQ Topic 或消息头区分 |
| 通信协议 | 进程内 async generator | RocketMQ 请求 + Redis Stream 响应 |
| 状态存储 | 进程内 + MySQL checkpointer | MySQL checkpointer(Agent Worker 管理) |
| 扩展方式 | 整体扩容 API 进程 | Agent Worker 独立扩容 |
| 流式输出 | async generator → Queue → SSE | Redis Stream → xread → SSE |
类模型设计
新增模块结构
text
backend/
├── agent/ # 不变:LangGraph Agent 编排层
│ ├── chat/
│ ├── image/
│ ├── video/
│ └── shared/
├── agent_worker/ # 新增:Agent Worker 独立入口
│ ├── __init__.py
│ ├── __main__.py # python -m agent_worker 入口
│ ├── bootstrap.py # Worker 启动基础设施(复用 worker/bootstrap 模式)
│ ├── config.py # Agent Worker 专属配置(RocketMQ topic、ConsumerGroup)
│ ├── consumer.py # RocketMQ 消费循环 + 消息分发
│ ├── handlers/ # 三大 Agent 的消息处理器
│ │ ├── __init__.py
│ │ ├── base_handler.py # AgentHandler 抽象基类
│ │ ├── chat_handler.py # ChatAgent 处理器
│ │ ├── image_handler.py # ImageAgent 处理器
│ │ └── video_handler.py # VideoAgent 处理器
│ ├── event_writer.py # Agent 事件 → Redis Stream 写入器
│ ├── adapters/ # 适配层:为 agent/ 提供独立依赖
│ │ ├── __init__.py
│ │ ├── checkpointer.py # 独立 checkpointer 管理(从 services/ 剥离)
│ │ ├── web_search.py # 独立 web_search 服务(从 services/ 剥离)
│ │ ├── agent_config.py # 独立推理配置读取(从 services/ 剥离)
│ │ └── providers.py # 独立 Provider 初始化
│ └── healthcheck.py # Docker 健康检查端点
├── Dockerfile.agent # 新增:Agent Worker Docker 镜像
└── ...核心类设计
依赖注入方式
各 Handler 通过构造函数注入 EventWriter,不直接持有 Redis 引用。 Agent 的 checkpointer、web_search、agent_config 等依赖通过 adapters/ 模块提供独立实现,不依赖 services/ 层。
API 侧改动
API 侧的 ChatStreamService、ImageStreamService、VideoStreamService 改造要点:
- 不再直接调用 Agent(移除
from agent.chat import get_chat_agent等导入) - 通过
AgentRequestPublisher发布 RocketMQ 消息 - 后台任务改为从 Redis Stream 读取事件(替代原来的 Agent async generator)
状态与数据流
端到端请求流(以 Chat 为例)
RocketMQ 消息体设计
Topic:SILIANG_AGENT_REQUEST_TOPIC(新增,NORMAL 类型)
消息体:
json
{
"request_id": "uuid-v4",
"agent_type": "chat | image | video",
"message_id": 12345,
"user_id": 100,
"thread_id": "thread-abc",
"conversation_id": "conv-xyz",
"prompt": "帮我生成一张猫的图片",
"model": "gpt-image-2",
"provider": {
"provider": "openrouter",
"api_keys": {"openrouter": "sk-xxx"}
},
"skill": {
"enabled": true,
"skill_source": "platform",
"platform_skill_id": 5,
"current_skill_id": 12,
"skill_project_dir": "/data/skills/xxx",
"skill_prompt_context": "You are ..."
},
"input": {
"web_search_enabled": false,
"images": [],
"temperature": 0.7,
"aspect_ratio": "1:1",
"reference_images": [],
"duration": 5,
"video_mode": "default"
},
"langsmith": {
"run_name": "chat-stream-abc",
"tags": ["chat", "stream", "gpt-5"],
"metadata": {"user_id": 100}
},
"timestamp": "2026-06-05T10:00:00Z"
}Redis Stream 事件格式
Key:chat:stream:{message_id}(复用现有命名约定)
事件字段(与现有 resume/stream_service.py 的 append_event 格式一致):
json
{
"type": "text-delta | progress | data-imageGeneration | data-videoGeneration | completed | error | aborted",
"payload": "{\"text\": \"你好\"}",
"created_at": "2026-06-05T18:00:00+08:00"
}错误边界
| 场景 | 处理方式 |
|---|---|
| RocketMQ 消息发送失败 | API 侧返回 503,前端提示重试 |
| Agent Worker 处理超时 | Worker 内置超时(默认 300s),超时后写 error 终态事件到 Redis Stream |
| Agent Worker 崩溃 | 租约过期(60s),API 侧检测到租约失效后写 failed 终态事件 |
| Redis Stream 写入失败 | Agent Worker 记录错误日志 + ACK 消息(避免重复消费),API 侧租约超时兜底 |
| Agent 内部异常 | Handler 捕获异常,写 error 事件到 Redis Stream,ACK 消息 |
Mermaid 图示
模块依赖图(拆分后)
Agent Worker 内部流程
API 侧改造流程
日志设计
Agent Worker 日志
| 业务事件 | 日志级别 | JSON 字段 |
|---|---|---|
| Worker 启动/停止 | INFO | source, pid, consumer_group |
| 消息消费成功 | INFO | request_id, agent_type, message_id, user_id |
| 消息处理开始 | INFO | request_id, agent_type, model, prompt_prefix |
| 消息处理完成 | INFO | request_id, duration_ms, event_count |
| 消息处理异常 | ERROR | request_id, error, trace_id |
| Redis Stream 写入失败 | ERROR | request_id, message_id, redis_error |
| Checkpointer 连接异常 | ERROR | thread_id, error |
trace_id 传递
- API 侧生成
request_id(UUID v4),作为 RocketMQ 消息的keys字段传递。 - Agent Worker 将
request_id注入日志上下文(LogContextFilter),贯穿整个请求链路。 - 消息
keys字段格式:agent_request:{request_id}。
敏感信息保护
api_keys不写入日志(仅记录provider名称)。prompt仅记录前 50 字符(prompt_prefix)。- Redis Stream 中的事件 payload 不包含
api_keys。
实施步骤
阶段一:Agent Worker 基础骨架(可验证)
- 创建
backend/agent_worker/目录结构。 - 实现
AgentHandler抽象基类和EventWriter。 - 实现
AgentWorkerConsumer和消息分发逻辑。 - 实现
adapters/checkpointer.py(从services/checkpointer.py提取)。 - 实现
ChatHandler作为第一个 handler。 - 实现
__main__.py启动入口和bootstrap.py。 - 创建
Dockerfile.agent。 - 验证:Agent Worker 能消费 RocketMQ 消息、执行 ChatAgent、写入 Redis Stream。
阶段二:API 侧改造
- 实现
AgentRequestPublisher(复用rocketmq_integration/producer.py)。 - 改造
ChatStreamService._run_resumable_generation():- 移除
from agent.chat import get_chat_agent。 - 改为发布 RocketMQ 消息 → 后台任务从 Redis Stream 读取事件。
- 移除
- 调整 RocketMQ 配置(新增 Topic
SILIANG_AGENT_REQUEST_TOPIC)。 - 验证:前端发起 Chat 请求 → SSE 流式输出正常。
阶段三:Image / Video Handler
- 实现
ImageHandler和VideoHandler。 - 实现
adapters/web_search.py、adapters/agent_config.py、adapters/providers.py。 - 改造
ImageStreamService和VideoStreamService。 - 验证:图片生成和视频生成端到端正常。
阶段四:容器编排与部署
- 更新
docker/docker-compose.dev.yml:新增siliang-agent服务。 - 更新
docker/docker-compose.prod.yml:新增siliang-agent服务。 - 更新
backend/Makefile:新增agent/agent-build命令。 - 更新 RocketMQ Broker 配置:新增 Agent Request Topic。
- 验证:Docker Compose 一键启动完整环境。
阶段五:清理与文档
- 从
backend/services/core/的 stream service 中移除对agent/的直接导入。 agent/shared/skill_runtime_resolver.py不再被 Agent Worker 使用(Skill 由 API 侧解析后序列化传递),移除 Agent Worker 侧的依赖。- 更新
backend/AGENTS.md架构图。 - 更新
docs/相关文档。 - 补充测试。
测试与验收
自动化测试
| 测试类型 | 覆盖范围 |
|---|---|
| 单元测试 | AgentHandler 消息解析、EventWriter Redis 写入、AgentRequestPublisher 消息构建 |
| 集成测试 | Agent Worker 端到端:消费消息 → 执行 Agent → 写 Redis Stream → API 读取 |
| 合约测试 | RocketMQ 消息体格式前后兼容 |
| 续接测试 | 断线续接:Worker 写 Redis Stream → API xread → SSE 续接 |
手动验证路径
- Chat:前端发起聊天 → SSE 流式文本 → 中间断线 → 自动续接 → 消息完整保存。
- Image:前端发起生图 → SSE 推理过程 → 图片 URL 返回 → 会话历史记录。
- Video:前端发起视频生成 → SSE 推理过程 → 视频任务提交 → 轮询状态。
- 并发:多用户同时请求 → Worker 正确分发 → 无消息丢失。
- 容错:Worker 崩溃 → 租约过期 → API 侧正确返回错误。
边界用例
| 用例 | 预期行为 |
|---|---|
| Agent 处理超时(>300s) | Worker 写 error 事件,API 返回错误 |
| Redis Stream 满载(2000 条上限) | 自动淘汰旧事件,续接从最新事件开始 |
| RocketMQ 消息重复投递 | Handler 幂等处理(基于 request_id 去重或状态检查) |
| Worker 滚动重启 | 旧消息 ACK 后优雅退出,新 Worker 接管消费 |
| API 实例扩缩容 | 无状态,不影响进行中的 Agent 处理 |
风险与回滚
兼容性风险
| 风险 | 影响 | 缓解措施 |
|---|---|---|
| SSE 延迟增加 | 新增 RocketMQ → Worker → Redis Stream 链路,端到端延迟增加 50-200ms | 前端已有 loading 状态,用户感知有限 |
| Redis Stream 内存占用 | Agent 事件存储在 Redis 中(每条消息最多 2000 事件) | 已有 STREAM_MAX_LEN=2000 上限 + TTL 自动清理 |
| 消息格式不兼容 | API 和 Worker 版本不一致导致消息解析失败 | 消息体包含 version 字段,向下兼容 |
迁移风险
| 风险 | 影响 | 缓解措施 |
|---|---|---|
| 进行中的流请求中断 | 部署时正在处理的 Agent 请求会断开 | 滚动部署:先启动新 Worker → 切换 API 发布 → 关闭旧 Worker |
| Checkpointer 共享 | API 侧旧代码和 Worker 新代码同时访问同一 MySQL checkpointer | MySQL checkpointer 本身支持并发访问,无冲突 |
回滚方案
- API 侧的
StreamService保留条件分支:环境变量AGENT_SERVICE_MODE=embedded | remote。 embedded模式走原有进程内调用路径,remote模式走 RocketMQ + Redis Stream。- 回滚时将
AGENT_SERVICE_MODE切回embedded,重启 API 服务即可。