Skip to content

Agent 服务拆分:将 LangGraph Agent 独立为无状态 API + Agent Worker 架构 ​

背景与目标 ​

要解决的问题 ​

当前 backend/agent/ 中的三大核心 Agent(ChatAgent、ImageAgent、VideoAgent)与后端 FastAPI 服务运行在同一进程中。这导致:

  1. API 服务有状态:Agent 通过 LangGraph checkpointer 和上下文管理持有会话状态,API 服务无法真正无状态水平扩展。
  2. 资源竞争:LLM 推理的长时间连接和内存占用影响 API 服务的响应延迟和稳定性。
  3. 独立扩展受限:无法单独为 Agent 推理层增加计算资源(GPU、内存),必须整体扩容。
  4. 故障隔离弱:Agent 内部异常(如 LLM 超时、工具调用失败)可能影响 API 服务的其他请求。

可感知结果 ​

  • 后端 API 服务变为无状态服务,可自由水平扩展。
  • Agent 推理作为独立 Worker 进程运行,可独立部署、扩容、回滚。
  • 前端用户无感知:SSE 流式体验保持一致。
  • 新架构复用现有 RocketMQ + Redis Stream 基础设施,与 RAG/Review Worker 模式统一。

现状分析 ​

当前架构 ​

Agent 层的外部依赖 ​

依赖来源使用位置拆分后归属
services.checkpointerservices/checkpointer.py三个 Agent 的 *_stream() 方法Agent Worker 自带(MySQL 连接)
core.databasecore/database.pyImage/Video Agent 的 _build_runtime_config()Agent Worker 自带(读取推理配置)
services.agent_config_serviceservices/agent_config_service.pyImage/Video Agent 的 _build_runtime_config()Agent Worker 自带
services.web_search_serviceservices/web_search_service.pyagent/shared/web_search_tool.pyAgent Worker 自带
services.platform_skill_serviceservices/platform_skill_service.pyagent/shared/skill_runtime_resolver.py保留在 API 服务(由 Service 层解析后序列化传递)
services.review.skill.serviceservices/review/skill/service.pyagent/shared/skill_runtime_resolver.py保留在 API 服务(同上)
core.user_info.UserInfocore/user_info.py三个 Agent 构建 LangSmith configAgent Worker 自带(纯数据类)
core.process_poolcore/process_pool.pyVideo Agent 的 prompt_optimizerAgent Worker 自带
Provider 初始化core/providers.pyImage/Video tools 中调用 generate_image 等Agent Worker 自带

可复用的现有能力 ​

  1. RocketMQ 集成:rocketmq_integration/ 已封装 Producer/Consumer,支持自动重连和延迟初始化。
  2. Redis Stream 续接:services/core/resume/ 已实现完整的 Stream 读写、并发控制、租约心跳机制。
  3. Worker 基础设施:worker/bootstrap.py 提供日志、DB 连接池、Redis、信号处理等共享启动逻辑。
  4. 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 进程内
StreamBridgeAPI 侧桥接 Redis Stream → SSE客户端断线后可续接API 侧(已存在)
CheckpointerSessionLangGraph checkpoint 的数据库连接会话并发连接数受限(当前上限 15)Agent Worker 进程内

协作关系 ​

变化点 ​

变化维度当前拆分后
Agent 类型Chat / Image / Video 三种同,通过 RocketMQ Topic 或消息头区分
通信协议进程内 async generatorRocketMQ 请求 + Redis Stream 响应
状态存储进程内 + MySQL checkpointerMySQL checkpointer(Agent Worker 管理)
扩展方式整体扩容 API 进程Agent Worker 独立扩容
流式输出async generator → Queue → SSERedis 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 启动/停止INFOsource, pid, consumer_group
消息消费成功INFOrequest_id, agent_type, message_id, user_id
消息处理开始INFOrequest_id, agent_type, model, prompt_prefix
消息处理完成INFOrequest_id, duration_ms, event_count
消息处理异常ERRORrequest_id, error, trace_id
Redis Stream 写入失败ERRORrequest_id, message_id, redis_error
Checkpointer 连接异常ERRORthread_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 基础骨架(可验证) ​

  1. 创建 backend/agent_worker/ 目录结构。
  2. 实现 AgentHandler 抽象基类和 EventWriter。
  3. 实现 AgentWorkerConsumer 和消息分发逻辑。
  4. 实现 adapters/checkpointer.py(从 services/checkpointer.py 提取)。
  5. 实现 ChatHandler 作为第一个 handler。
  6. 实现 __main__.py 启动入口和 bootstrap.py。
  7. 创建 Dockerfile.agent。
  8. 验证:Agent Worker 能消费 RocketMQ 消息、执行 ChatAgent、写入 Redis Stream。

阶段二:API 侧改造 ​

  1. 实现 AgentRequestPublisher(复用 rocketmq_integration/producer.py)。
  2. 改造 ChatStreamService._run_resumable_generation():
    • 移除 from agent.chat import get_chat_agent。
    • 改为发布 RocketMQ 消息 → 后台任务从 Redis Stream 读取事件。
  3. 调整 RocketMQ 配置(新增 Topic SILIANG_AGENT_REQUEST_TOPIC)。
  4. 验证:前端发起 Chat 请求 → SSE 流式输出正常。

阶段三:Image / Video Handler ​

  1. 实现 ImageHandler 和 VideoHandler。
  2. 实现 adapters/web_search.py、adapters/agent_config.py、adapters/providers.py。
  3. 改造 ImageStreamService 和 VideoStreamService。
  4. 验证:图片生成和视频生成端到端正常。

阶段四:容器编排与部署 ​

  1. 更新 docker/docker-compose.dev.yml:新增 siliang-agent 服务。
  2. 更新 docker/docker-compose.prod.yml:新增 siliang-agent 服务。
  3. 更新 backend/Makefile:新增 agent / agent-build 命令。
  4. 更新 RocketMQ Broker 配置:新增 Agent Request Topic。
  5. 验证:Docker Compose 一键启动完整环境。

阶段五:清理与文档 ​

  1. 从 backend/services/core/ 的 stream service 中移除对 agent/ 的直接导入。
  2. agent/shared/skill_runtime_resolver.py 不再被 Agent Worker 使用(Skill 由 API 侧解析后序列化传递),移除 Agent Worker 侧的依赖。
  3. 更新 backend/AGENTS.md 架构图。
  4. 更新 docs/ 相关文档。
  5. 补充测试。

测试与验收 ​

自动化测试 ​

测试类型覆盖范围
单元测试AgentHandler 消息解析、EventWriter Redis 写入、AgentRequestPublisher 消息构建
集成测试Agent Worker 端到端:消费消息 → 执行 Agent → 写 Redis Stream → API 读取
合约测试RocketMQ 消息体格式前后兼容
续接测试断线续接:Worker 写 Redis Stream → API xread → SSE 续接

手动验证路径 ​

  1. Chat:前端发起聊天 → SSE 流式文本 → 中间断线 → 自动续接 → 消息完整保存。
  2. Image:前端发起生图 → SSE 推理过程 → 图片 URL 返回 → 会话历史记录。
  3. Video:前端发起视频生成 → SSE 推理过程 → 视频任务提交 → 轮询状态。
  4. 并发:多用户同时请求 → Worker 正确分发 → 无消息丢失。
  5. 容错: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 checkpointerMySQL checkpointer 本身支持并发访问,无冲突

回滚方案 ​

  1. API 侧的 StreamService 保留条件分支:环境变量 AGENT_SERVICE_MODE=embedded | remote。
  2. embedded 模式走原有进程内调用路径,remote 模式走 RocketMQ + Redis Stream。
  3. 回滚时将 AGENT_SERVICE_MODE 切回 embedded,重启 API 服务即可。