# 队列服务 **本文引用的文件** - [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)