# 队列服务
**本文引用的文件**
- [queue.service.ts](file://server/src/services/queue.service.ts)
- [memory-queue.ts](file://server/src/services/memory-queue.ts)
- [redis.service.ts](file://server/src/services/redis.service.ts)
- [book-queue.processor.ts](file://server/src/modules/book-generator/book-queue.processor.ts)
- [fault-tolerance.ts](file://server/src/modules/book-generator/fault-tolerance.ts)
- [index.ts](file://server/src/modules/book-generator/index.ts)
## 目录
1. [简介](#简介)
2. [项目结构](#项目结构)
3. [核心组件](#核心组件)
4. [架构总览](#架构总览)
5. [详细组件分析](#详细组件分析)
6. [依赖关系分析](#依赖关系分析)
7. [性能考量](#性能考量)
8. [故障排查指南](#故障排查指南)
9. [结论](#结论)
10. [附录](#附录)
## 简介
本文件面向AI有声书生成平台的“队列服务”,系统化阐述基于Bull的异步任务处理机制,覆盖以下主题:
- 任务队列的创建、调度、并发控制与优先级管理
- 失败重试与超时控制的职责划分与实现位置
- 内存队列的实现原理与回退策略
- 队列处理器的工作流程与任务状态跟踪
- 队列配置项、性能调优参数与故障恢复机制
- 在AI书籍生成、TTS语音合成、视频生成等核心业务中的应用示例与最佳实践
## 项目结构
队列服务位于后端服务层,围绕“单一职责”原则设计:队列仅负责排队与并发控制;失败重试与超时管理由“容错层”承担;业务逻辑由领域层实现。
```mermaid
graph TB
subgraph "服务层"
QS["QueueService
队列服务"]
MQ["MemoryQueue
内存队列"]
RS["RedisService
Redis服务"]
end
subgraph "业务模块"
BG["book-queue.processor.ts
书籍生成队列处理器"]
FT["fault-tolerance.ts
容错层"]
LG["index.ts
LangGraph生成器"]
end
QS --> RS
QS --> MQ
BG --> QS
BG --> MQ
BG --> LG
FT --> QS
FT --> BG
```
图表来源
- [queue.service.ts:18-347](file://server/src/services/queue.service.ts#L18-L347)
- [memory-queue.ts:17-119](file://server/src/services/memory-queue.ts#L17-L119)
- [redis.service.ts:3-274](file://server/src/services/redis.service.ts#L3-L274)
- [book-queue.processor.ts:7-83](file://server/src/modules/book-generator/book-queue.processor.ts#L7-L83)
- [fault-tolerance.ts:11-387](file://server/src/modules/book-generator/fault-tolerance.ts#L11-L387)
- [index.ts:74-91](file://server/src/modules/book-generator/index.ts#L74-L91)
章节来源
- [queue.service.ts:18-347](file://server/src/services/queue.service.ts#L18-L347)
- [memory-queue.ts:17-119](file://server/src/services/memory-queue.ts#L17-L119)
- [redis.service.ts:3-274](file://server/src/services/redis.service.ts#L3-L274)
- [book-queue.processor.ts:7-83](file://server/src/modules/book-generator/book-queue.processor.ts#L7-L83)
- [fault-tolerance.ts:11-387](file://server/src/modules/book-generator/fault-tolerance.ts#L11-L387)
- [index.ts:74-91](file://server/src/modules/book-generator/index.ts#L74-L91)
## 核心组件
- 队列服务(QueueService)
- 负责:队列实例化、任务入队、状态查询、进度更新、统计与运维操作
- 特性:Redis队列为主,内存队列为回退;统一的超时与去重策略;事件监听与可用性检测
- 内存队列(MemoryQueue)
- 负责:在Redis不可用时提供基础队列能力(并发、等待队列、处理循环)
- 特性:基于内存Map与数组的轻量实现,支持并发槽位与简单事件钩子
- Redis服务(RedisService)
- 负责:Redis连接生命周期、可用性检测、基础KV操作
- 特性:自动重连策略、错误隔离与连接状态维护
- 书籍生成队列处理器(book-queue.processor.ts)
- 负责:注册Redis或内存队列处理器,恢复中断任务,推进生成阶段
- 特性:并发控制(默认3)、进度事件监听、失败状态回写
- 容错层(fault-tolerance.ts)
- 负责:AI调用重试(指数退避)、节点超时控制、进度监控与自动恢复
- 特性:配置化的重试次数与超时阈值;用户通知与失败日志记录
章节来源
- [queue.service.ts:48-347](file://server/src/services/queue.service.ts#L48-L347)
- [memory-queue.ts:17-119](file://server/src/services/memory-queue.ts#L17-L119)
- [redis.service.ts:3-274](file://server/src/services/redis.service.ts#L3-L274)
- [book-queue.processor.ts:16-83](file://server/src/modules/book-generator/book-queue.processor.ts#L16-L83)
- [fault-tolerance.ts:17-387](file://server/src/modules/book-generator/fault-tolerance.ts#L17-L387)
## 架构总览
队列服务采用“主从双通道”架构:优先使用Redis队列承载高并发与持久化;当Redis不可用时,自动切换至内存队列以保证系统可用性。书籍生成等长耗时任务通过处理器按并发上限执行,并结合容错层实现稳健的失败恢复与进度监控。
```mermaid
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](file://server/src/services/queue.service.ts#L131-L160)
- [redis.service.ts:43-45](file://server/src/services/redis.service.ts#L43-L45)
- [memory-queue.ts:26-46](file://server/src/services/memory-queue.ts#L26-L46)
- [book-queue.processor.ts:52-82](file://server/src/modules/book-generator/book-queue.processor.ts#L52-L82)
## 详细组件分析
### 队列服务(QueueService)
- 设计原则
- 单一职责:仅负责排队与并发控制
- 失败重试与超时管理交由容错层与AI服务层处理
- 业务逻辑由领域层实现
- 关键能力
- 任务入队:支持不同类型任务(音频、视频、书籍、邮件),带超时与去重配置
- 状态查询:通过Bull Job.getState映射为统一状态枚举
- 进度更新:支持向Bull Job写入进度并触发本地回调
- 统计与运维:等待/活动/完成/失败/延迟计数;清空、暂停、恢复、关闭
- 可用性与回退:Redis不可用时自动切换内存队列
- 配置要点
- 默认完成/失败记录清理数量
- 队列空闲挂起检测与最大挂起次数
- Redis连接参数与重试策略
- 适用场景
- 音频生成:较长超时(5分钟)
- 视频生成:较长超时(10分钟)
- 书籍生成:极长超时(2小时),配合容错层进行稳健恢复
```mermaid
classDiagram
class QueueService {
+isQueueAvailable() boolean
+addTask(queueType, data, options) Promise
+addAudioGenerationTask(data) Promise
+addVideoGenerationTask(data) Promise
+addBookGenerationTask(data) Promise
+getTaskStatus(queueType, jobId) Promise
+updateProgress(queueType, jobId, progress, data) Promise
+onProgress(jobId, callback) void
+getQueueStats(queueType) Promise
+clearQueue(queueType) Promise
+pauseQueue(queueType) Promise
+resumeQueue(queueType) Promise
+closeAll() Promise
}
class RedisService {
+isAvailable() boolean
+get(key) Promise
+set(key, value, ttl) Promise
+del(key) Promise
+testConnection() Promise
}
class MemoryQueue {
+add(name, data) Promise
+process(concurrency, handler) void
+hasPendingJobs() boolean
+close() Promise
}
QueueService --> RedisService : "依赖"
QueueService --> MemoryQueue : "回退"
```
图表来源
- [queue.service.ts:48-347](file://server/src/services/queue.service.ts#L48-L347)
- [redis.service.ts:3-274](file://server/src/services/redis.service.ts#L3-L274)
- [memory-queue.ts:17-119](file://server/src/services/memory-queue.ts#L17-L119)
章节来源
- [queue.service.ts:48-347](file://server/src/services/queue.service.ts#L48-L347)
- [redis.service.ts:3-274](file://server/src/services/redis.service.ts#L3-L274)
- [memory-queue.ts:17-119](file://server/src/services/memory-queue.ts#L17-L119)
### 内存队列(MemoryQueue)
- 实现原理
- 使用Map存储任务元数据,数组维护等待队列,Set跟踪正在处理的任务
- 通过并发槽位控制同时处理的任务数,处理循环在空闲时延时,避免忙等
- 提供简单的事件钩子占位,便于扩展
- 使用场景
- Redis不可用时的临时回退
- 开发/测试环境或低并发场景
- 注意事项
- 数据驻留内存,进程重启即丢失
- 不具备持久化与跨进程共享能力
```mermaid
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](file://server/src/services/memory-queue.ts#L26-L99)
章节来源
- [memory-queue.ts:17-119](file://server/src/services/memory-queue.ts#L17-L119)
### 书籍生成队列处理器(book-queue.processor.ts)
- 处理器注册
- 优先使用Redis队列,设置并发数为3
- 若Redis不可用,则使用内存队列处理器
- 任务处理流程
- 从job.data提取参数,更新书籍状态为“大纲生成”
- 调用LangGraph生成器执行生成
- 成功则返回结果,失败则更新书籍状态为“失败”
- 中断任务恢复
- 启动时扫描数据库中处于“内容/音频/视频生成中”的书籍
- 将其重新加入队列,恢复生成流程
```mermaid
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](file://server/src/modules/book-generator/book-queue.processor.ts#L16-L43)
- [index.ts:74-91](file://server/src/modules/book-generator/index.ts#L74-L91)
章节来源
- [book-queue.processor.ts:16-83](file://server/src/modules/book-generator/book-queue.processor.ts#L16-L83)
- [index.ts:74-91](file://server/src/modules/book-generator/index.ts#L74-L91)
### 容错层(fault-tolerance.ts)
- 重试机制
- AI调用失败自动重试,指数退避,最大重试次数可配置
- 每次重试均记录到书籍错误信息,便于追踪
- 超时控制
- 节点级超时控制,超时后通知用户并决定后续策略
- 进度监控
- 定期检查书籍进度,若长时间无进展发出警告,必要时触发自动恢复
- 自动恢复
- 在限定次数内自动将任务重新入队,降低人工干预成本
```mermaid
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](file://server/src/modules/book-generator/fault-tolerance.ts#L68-L123)
- [fault-tolerance.ts:131-180](file://server/src/modules/book-generator/fault-tolerance.ts#L131-L180)
- [fault-tolerance.ts:188-261](file://server/src/modules/book-generator/fault-tolerance.ts#L188-L261)
- [fault-tolerance.ts:268-323](file://server/src/modules/book-generator/fault-tolerance.ts#L268-L323)
章节来源
- [fault-tolerance.ts:17-387](file://server/src/modules/book-generator/fault-tolerance.ts#L17-L387)
## 依赖关系分析
- 组件耦合
- QueueService与RedisService强耦合,与MemoryQueue弱耦合(回退)
- 书籍生成处理器依赖QueueService与MemoryQueue,同时依赖LangGraph生成器与bookStore
- 容错层与QueueService、书籍生成处理器形成松耦合协作
- 外部依赖
- Bull(Redis队列)、ioredis(Redis客户端)
- 循环依赖
- 未发现循环依赖迹象
```mermaid
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](file://server/src/services/queue.service.ts#L18-L20)
- [book-queue.processor.ts:7-11](file://server/src/modules/book-generator/book-queue.processor.ts#L7-L11)
- [fault-tolerance.ts:11-13](file://server/src/modules/book-generator/fault-tolerance.ts#L11-L13)
章节来源
- [queue.service.ts:18-20](file://server/src/services/queue.service.ts#L18-L20)
- [book-queue.processor.ts:7-11](file://server/src/modules/book-generator/book-queue.processor.ts#L7-L11)
- [fault-tolerance.ts:11-13](file://server/src/modules/book-generator/fault-tolerance.ts#L11-L13)
## 性能考量
- 并发控制
- 书籍生成处理器默认并发3,可根据服务器资源与AI服务吞吐调整
- 超时与去重
- 任务级别超时避免长时间占用资源;完成/失败记录清理减少Redis冗余
- Redis优化
- 合理设置stalledInterval与maxStalledCount,平衡任务恢复与资源消耗
- 内存队列
- 适合低并发场景;生产环境建议使用Redis队列以获得更好的扩展性与可靠性
## 故障排查指南
- Redis不可用
- 现象:队列服务不可用标志被置为false,任务入队回退到内存队列
- 处理:检查Redis连接参数与网络;使用RedisService的测试方法验证连通性
- 任务长时间无响应
- 现象:进度监控发出“长时间无响应”警告,自动恢复尝试
- 处理:检查AI服务可用性与超时配置;确认容错层日志与用户通知
- 任务失败
- 现象:队列事件监听输出失败日志,进度回调清理
- 处理:查看任务状态与失败原因;结合容错层重试与自动恢复策略
章节来源
- [queue.service.ts:97-110](file://server/src/services/queue.service.ts#L97-L110)
- [redis.service.ts:24-37](file://server/src/services/redis.service.ts#L24-L37)
- [fault-tolerance.ts:188-261](file://server/src/modules/book-generator/fault-tolerance.ts#L188-L261)
## 结论
队列服务通过“Redis主通道 + 内存回退”的设计,在保证高并发与可靠性的同时,兼顾了可用性与可维护性。配合容错层的重试、超时与自动恢复机制,能够有效应对AI服务波动与长耗时任务的不确定性。在AI书籍生成、TTS语音合成、视频生成等核心业务中,该架构提供了稳定、可观测、可扩展的任务处理能力。
## 附录
### 队列配置与调优清单
- Redis连接参数
- 主机、端口、密码、数据库索引
- 队列行为参数
- 默认完成/失败记录清理数量
- 僵尸任务检测间隔与最大挂起次数
- 任务超时
- 音频生成:约5分钟
- 视频生成:约10分钟
- 书籍生成:约2小时
- 并发控制
- 书籍生成处理器默认并发3,可根据资源与SLA调整
章节来源
- [queue.service.ts:79-95](file://server/src/services/queue.service.ts#L79-L95)
- [queue.service.ts:166-190](file://server/src/services/queue.service.ts#L166-L190)
- [book-queue.processor.ts:56](file://server/src/modules/book-generator/book-queue.processor.ts#L56)
### 实际使用示例(路径指引)
- 创建书籍生成任务
- 调用:[addBookGenerationTask:186-190](file://server/src/services/queue.service.ts#L186-L190)
- 处理器:[initBookGenerationQueue:48-83](file://server/src/modules/book-generator/book-queue.processor.ts#L48-L83)
- 监控任务状态
- 查询:[getTaskStatus:197-237](file://server/src/services/queue.service.ts#L197-L237)
- 统计:[getQueueStats:290-300](file://server/src/services/queue.service.ts#L290-L300)
- 更新任务进度
- 更新:[updateProgress:246-267](file://server/src/services/queue.service.ts#L246-L267)
- 回调:[onProgress:274-276](file://server/src/services/queue.service.ts#L274-L276)
- 处理队列异常
- Redis不可用回退:[addTask回退逻辑:136-160](file://server/src/services/queue.service.ts#L136-L160)
- 容错重试与恢复:[fault-tolerance:68-323](file://server/src/modules/book-generator/fault-tolerance.ts#L68-L323)
- 恢复中断任务
- 恢复:[resumeInterruptedTasks:88-124](file://server/src/modules/book-generator/book-queue.processor.ts#L88-L124)