队列服务.md 17 KB

队列服务

本文引用的文件

  • queue.service.ts
  • memory-queue.ts
  • redis.service.ts
  • book-queue.processor.ts
  • fault-tolerance.ts
  • index.ts

目录

  1. 简介
  2. 项目结构
  3. 核心组件
  4. 架构总览
  5. 详细组件分析
  6. 依赖关系分析
  7. 性能考量
  8. 故障排查指南
  9. 结论
  10. 附录

简介

本文件面向AI有声书生成平台的“队列服务”,系统化阐述基于Bull的异步任务处理机制,覆盖以下主题:

  • 任务队列的创建、调度、并发控制与优先级管理
  • 失败重试与超时控制的职责划分与实现位置
  • 内存队列的实现原理与回退策略
  • 队列处理器的工作流程与任务状态跟踪
  • 队列配置项、性能调优参数与故障恢复机制
  • 在AI书籍生成、TTS语音合成、视频生成等核心业务中的应用示例与最佳实践

项目结构

队列服务位于后端服务层,围绕“单一职责”原则设计:队列仅负责排队与并发控制;失败重试与超时管理由“容错层”承担;业务逻辑由领域层实现。

graph TB
subgraph "服务层"
QS["QueueService<br/>队列服务"]
MQ["MemoryQueue<br/>内存队列"]
RS["RedisService<br/>Redis服务"]
end
subgraph "业务模块"
BG["book-queue.processor.ts<br/>书籍生成队列处理器"]
FT["fault-tolerance.ts<br/>容错层"]
LG["index.ts<br/>LangGraph生成器"]
end
QS --> RS
QS --> MQ
BG --> QS
BG --> MQ
BG --> LG
FT --> QS
FT --> BG

图表来源

  • queue.service.ts:18-347
  • memory-queue.ts:17-119
  • redis.service.ts:3-274
  • book-queue.processor.ts:7-83
  • fault-tolerance.ts:11-387
  • index.ts:74-91

章节来源

  • queue.service.ts:18-347
  • memory-queue.ts:17-119
  • redis.service.ts:3-274
  • book-queue.processor.ts:7-83
  • fault-tolerance.ts:11-387
  • index.ts:74-91

核心组件

  • 队列服务(QueueService)
    • 负责:队列实例化、任务入队、状态查询、进度更新、统计与运维操作
    • 特性:Redis队列为主,内存队列为回退;统一的超时与去重策略;事件监听与可用性检测
  • 内存队列(MemoryQueue)
    • 负责:在Redis不可用时提供基础队列能力(并发、等待队列、处理循环)
    • 特性:基于内存Map与数组的轻量实现,支持并发槽位与简单事件钩子
  • Redis服务(RedisService)
    • 负责:Redis连接生命周期、可用性检测、基础KV操作
    • 特性:自动重连策略、错误隔离与连接状态维护
  • 书籍生成队列处理器(book-queue.processor.ts)
    • 负责:注册Redis或内存队列处理器,恢复中断任务,推进生成阶段
    • 特性:并发控制(默认3)、进度事件监听、失败状态回写
  • 容错层(fault-tolerance.ts)
    • 负责:AI调用重试(指数退避)、节点超时控制、进度监控与自动恢复
    • 特性:配置化的重试次数与超时阈值;用户通知与失败日志记录

章节来源

  • queue.service.ts:48-347
  • memory-queue.ts:17-119
  • redis.service.ts:3-274
  • book-queue.processor.ts:16-83
  • fault-tolerance.ts:17-387

架构总览

队列服务采用“主从双通道”架构:优先使用Redis队列承载高并发与持久化;当Redis不可用时,自动切换至内存队列以保证系统可用性。书籍生成等长耗时任务通过处理器按并发上限执行,并结合容错层实现稳健的失败恢复与进度监控。

sequenceDiagram
participant C as "调用方"
participant QS as "QueueService"
participant RS as "RedisService"
participant RQ as "Redis队列(Bull)"
participant MQ as "MemoryQueue"
participant BG as "书籍生成处理器"
C->>QS : "addBookGenerationTask(data)"
QS->>RS : "isAvailable()"
alt Redis可用
QS->>RQ : "queue.add(data, options)"
RQ-->>QS : "jobId"
QS-->>C : "jobId"
RQ-->>BG : "process(3, handler)"
BG-->>RQ : "进度事件(progress)"
RQ-->>BG : "完成(completed)/失败(failed)"
else Redis不可用
QS->>MQ : "memoryQueue.add(type, data)"
MQ-->>QS : "jobId"
QS-->>C : "jobId"
MQ-->>BG : "process(handler)"
end

图表来源

  • queue.service.ts:131-160
  • redis.service.ts:43-45
  • memory-queue.ts:26-46
  • book-queue.processor.ts:52-82

详细组件分析

队列服务(QueueService)

  • 设计原则
    • 单一职责:仅负责排队与并发控制
    • 失败重试与超时管理交由容错层与AI服务层处理
    • 业务逻辑由领域层实现
  • 关键能力
    • 任务入队:支持不同类型任务(音频、视频、书籍、邮件),带超时与去重配置
    • 状态查询:通过Bull Job.getState映射为统一状态枚举
    • 进度更新:支持向Bull Job写入进度并触发本地回调
    • 统计与运维:等待/活动/完成/失败/延迟计数;清空、暂停、恢复、关闭
    • 可用性与回退:Redis不可用时自动切换内存队列
  • 配置要点
    • 默认完成/失败记录清理数量
    • 队列空闲挂起检测与最大挂起次数
    • Redis连接参数与重试策略
  • 适用场景

    • 音频生成:较长超时(5分钟)
    • 视频生成:较长超时(10分钟)
    • 书籍生成:极长超时(2小时),配合容错层进行稳健恢复

      classDiagram
      class QueueService {
      +isQueueAvailable() boolean
      +addTask(queueType, data, options) Promise<string|null>
      +addAudioGenerationTask(data) Promise<string|null>
      +addVideoGenerationTask(data) Promise<string|null>
      +addBookGenerationTask(data) Promise<string|null>
      +getTaskStatus(queueType, jobId) Promise<Status>
      +updateProgress(queueType, jobId, progress, data) Promise<void>
      +onProgress(jobId, callback) void
      +getQueueStats(queueType) Promise<counts>
      +clearQueue(queueType) Promise<void>
      +pauseQueue(queueType) Promise<void>
      +resumeQueue(queueType) Promise<void>
      +closeAll() Promise<void>
      }
      class RedisService {
      +isAvailable() boolean
      +get(key) Promise<string|null>
      +set(key, value, ttl) Promise<boolean>
      +del(key) Promise<boolean>
      +testConnection() Promise<boolean>
      }
      class MemoryQueue {
      +add(name, data) Promise<string>
      +process(concurrency, handler) void
      +hasPendingJobs() boolean
      +close() Promise<void>
      }
      QueueService --> RedisService : "依赖"
      QueueService --> MemoryQueue : "回退"
      

图表来源

  • queue.service.ts:48-347
  • redis.service.ts:3-274
  • memory-queue.ts:17-119

章节来源

  • queue.service.ts:48-347
  • redis.service.ts:3-274
  • memory-queue.ts:17-119

内存队列(MemoryQueue)

  • 实现原理
    • 使用Map存储任务元数据,数组维护等待队列,Set跟踪正在处理的任务
    • 通过并发槽位控制同时处理的任务数,处理循环在空闲时延时,避免忙等
    • 提供简单的事件钩子占位,便于扩展
  • 使用场景
    • Redis不可用时的临时回退
    • 开发/测试环境或低并发场景
  • 注意事项

    • 数据驻留内存,进程重启即丢失
    • 不具备持久化与跨进程共享能力

      flowchart TD
      Start(["开始"]) --> Add["add(name, data)"]
      Add --> Push["加入等待队列"]
      Push --> CheckHandler{"已有处理器?"}
      CheckHandler --> |是| StartProc["startProcessing()"]
      CheckHandler --> |否| Wait["等待处理器注册"]
      StartProc --> Loop["processLoop()"]
      Loop --> HasSlot{"有空闲槽位?"}
      HasSlot --> |否| Sleep["setTimeout(100ms)"] --> Loop
      HasSlot --> |是| Pop["弹出一个任务"]
      Pop --> MarkActive["标记为active"]
      MarkActive --> TryExec["调用handler(job)"]
      TryExec --> Ok{"执行成功?"}
      Ok --> |是| MarkDone["标记为completed"]
      Ok --> |否| MarkFail["标记为failed"]
      MarkDone --> Next["setImmediate继续循环"] --> Loop
      MarkFail --> Next --> Loop
      

图表来源

  • memory-queue.ts:26-99

章节来源

  • memory-queue.ts:17-119

书籍生成队列处理器(book-queue.processor.ts)

  • 处理器注册
    • 优先使用Redis队列,设置并发数为3
    • 若Redis不可用,则使用内存队列处理器
  • 任务处理流程
    • 从job.data提取参数,更新书籍状态为“大纲生成”
    • 调用LangGraph生成器执行生成
    • 成功则返回结果,失败则更新书籍状态为“失败”
  • 中断任务恢复

    • 启动时扫描数据库中处于“内容/音频/视频生成中”的书籍
    • 将其重新加入队列,恢复生成流程

      sequenceDiagram
      participant BG as "书籍生成处理器"
      participant Store as "bookStore"
      participant LG as "LangGraph生成器"
      participant DB as "Prisma"
      BG->>Store : "update(bookId, {genStage : 'outlining', progress : 0})"
      BG->>LG : "generate(bookId, topic, bookScale, genLevel)"
      alt 成功
      LG-->>BG : "success"
      BG-->>Store : "update(bookId, {genStage : 'done'})"
      else 失败
      LG-->>BG : "error"
      BG-->>Store : "update(bookId, {genStage : 'failed', failedStage : 'outlining', errorMsg})"
      end
      

图表来源

  • book-queue.processor.ts:16-43
  • index.ts:74-91

章节来源

  • book-queue.processor.ts:16-83
  • index.ts:74-91

容错层(fault-tolerance.ts)

  • 重试机制
    • AI调用失败自动重试,指数退避,最大重试次数可配置
    • 每次重试均记录到书籍错误信息,便于追踪
  • 超时控制
    • 节点级超时控制,超时后通知用户并决定后续策略
  • 进度监控
    • 定期检查书籍进度,若长时间无进展发出警告,必要时触发自动恢复
  • 自动恢复

    • 在限定次数内自动将任务重新入队,降低人工干预成本

      flowchart TD
      Enter(["进入节点"]) --> Timeout["设定超时阈值"]
      Timeout --> CallLLM["callLLMWithRetry(...)"]
      CallLLM --> Retry{"重试中?"}
      Retry --> |是| Backoff["指数退避等待"] --> CallLLM
      Retry --> |否| Done["返回结果"]
      Done --> Monitor["startProgressMonitor(...)"]
      Monitor --> Idle{"长时间无进度?"}
      Idle --> |是| AutoRecovery["attemptAutoRecovery(...)"]
      Idle --> |否| Continue["继续生成"]
      AutoRecovery --> Requeue["重新入队"] --> Continue
      

图表来源

  • fault-tolerance.ts:68-123
  • fault-tolerance.ts:131-180
  • fault-tolerance.ts:188-261
  • fault-tolerance.ts:268-323

章节来源

  • fault-tolerance.ts:17-387

依赖关系分析

  • 组件耦合
    • QueueService与RedisService强耦合,与MemoryQueue弱耦合(回退)
    • 书籍生成处理器依赖QueueService与MemoryQueue,同时依赖LangGraph生成器与bookStore
    • 容错层与QueueService、书籍生成处理器形成松耦合协作
  • 外部依赖
    • Bull(Redis队列)、ioredis(Redis客户端)
  • 循环依赖

    • 未发现循环依赖迹象

      graph LR
      QS["QueueService"] --> RS["RedisService"]
      QS --> MQ["MemoryQueue"]
      BG["book-queue.processor.ts"] --> QS
      BG --> MQ
      BG --> LG["LangGraph生成器"]
      FT["fault-tolerance.ts"] --> QS
      FT --> BG
      

图表来源

  • queue.service.ts:18-20
  • book-queue.processor.ts:7-11
  • fault-tolerance.ts:11-13

章节来源

  • queue.service.ts:18-20
  • book-queue.processor.ts:7-11
  • fault-tolerance.ts:11-13

性能考量

  • 并发控制
    • 书籍生成处理器默认并发3,可根据服务器资源与AI服务吞吐调整
  • 超时与去重
    • 任务级别超时避免长时间占用资源;完成/失败记录清理减少Redis冗余
  • Redis优化
    • 合理设置stalledInterval与maxStalledCount,平衡任务恢复与资源消耗
  • 内存队列
    • 适合低并发场景;生产环境建议使用Redis队列以获得更好的扩展性与可靠性

故障排查指南

  • Redis不可用
    • 现象:队列服务不可用标志被置为false,任务入队回退到内存队列
    • 处理:检查Redis连接参数与网络;使用RedisService的测试方法验证连通性
  • 任务长时间无响应
    • 现象:进度监控发出“长时间无响应”警告,自动恢复尝试
    • 处理:检查AI服务可用性与超时配置;确认容错层日志与用户通知
  • 任务失败
    • 现象:队列事件监听输出失败日志,进度回调清理
    • 处理:查看任务状态与失败原因;结合容错层重试与自动恢复策略

章节来源

  • queue.service.ts:97-110
  • redis.service.ts:24-37
  • fault-tolerance.ts:188-261

结论

队列服务通过“Redis主通道 + 内存回退”的设计,在保证高并发与可靠性的同时,兼顾了可用性与可维护性。配合容错层的重试、超时与自动恢复机制,能够有效应对AI服务波动与长耗时任务的不确定性。在AI书籍生成、TTS语音合成、视频生成等核心业务中,该架构提供了稳定、可观测、可扩展的任务处理能力。

附录

队列配置与调优清单

  • Redis连接参数
    • 主机、端口、密码、数据库索引
  • 队列行为参数
    • 默认完成/失败记录清理数量
    • 僵尸任务检测间隔与最大挂起次数
  • 任务超时
    • 音频生成:约5分钟
    • 视频生成:约10分钟
    • 书籍生成:约2小时
  • 并发控制
    • 书籍生成处理器默认并发3,可根据资源与SLA调整

章节来源

  • queue.service.ts:79-95
  • queue.service.ts:166-190
  • book-queue.processor.ts:56

实际使用示例(路径指引)

  • 创建书籍生成任务
    • 调用:addBookGenerationTask:186-190
    • 处理器:initBookGenerationQueue:48-83
  • 监控任务状态
    • 查询:getTaskStatus:197-237
    • 统计:getQueueStats:290-300
  • 更新任务进度
    • 更新:updateProgress:246-267
    • 回调:onProgress:274-276
  • 处理队列异常
    • Redis不可用回退:addTask回退逻辑:136-160
    • 容错重试与恢复:fault-tolerance:68-323
  • 恢复中断任务
    • 恢复:resumeInterruptedTasks:88-124