本文引用的文件
本文件面向AI有声书生成平台的“队列服务”,系统化阐述基于Bull的异步任务处理机制,覆盖以下主题:
队列服务位于后端服务层,围绕“单一职责”原则设计:队列仅负责排队与并发控制;失败重试与超时管理由“容错层”承担;业务逻辑由领域层实现。
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
图表来源
章节来源
章节来源
队列服务采用“主从双通道”架构:优先使用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
图表来源
适用场景
书籍生成:极长超时(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 : "回退"
图表来源
章节来源
注意事项
不具备持久化与跨进程共享能力
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
图表来源
章节来源
中断任务恢复
将其重新加入队列,恢复生成流程
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
图表来源
章节来源
自动恢复
在限定次数内自动将任务重新入队,降低人工干预成本
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
图表来源
章节来源
循环依赖
未发现循环依赖迹象
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
图表来源
章节来源
章节来源
队列服务通过“Redis主通道 + 内存回退”的设计,在保证高并发与可靠性的同时,兼顾了可用性与可维护性。配合容错层的重试、超时与自动恢复机制,能够有效应对AI服务波动与长耗时任务的不确定性。在AI书籍生成、TTS语音合成、视频生成等核心业务中,该架构提供了稳定、可观测、可扩展的任务处理能力。
章节来源