批量处理机制.md 20 KB

批量处理机制

本文引用的文件

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

目录

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

简介

本文件系统性阐述本项目的批量处理机制,围绕“队列系统架构、任务调度算法、优先级与并发控制、资源限制与负载均衡、任务状态跟踪与进度监控、异常处理与超时控制、重试机制、内存队列与持久化队列协调、数据一致性、性能监控与吞吐优化、配置参数与扩展接口、故障恢复策略”等方面展开。文档以代码为依据,结合可视化图示帮助读者快速理解与落地实施。

项目结构

与批量处理相关的核心模块分布如下:

  • 服务层
    • 队列服务:统一管理 Redis/Bull 队列与内存队列回退
    • Redis 服务:提供连接状态与能力检测
    • WebSocket 服务:向客户端推送进度与事件
    • 内存队列:Redis 不可用时的本地队列实现
  • 业务层
    • 书籍生成队列处理器:绑定队列并发与事件处理
    • 容错层:AI 调用重试、节点超时、进度监控、自动恢复
    • 批量生成编排器:按步骤推进内容/音频/视频生成,并推送进度
  • 前端

    • WebSocket 工具:连接与事件订阅管理,支持 H5/小程序

      graph TB
      subgraph "服务层"
      QS["队列服务<br/>queue.service.ts"]
      RS["Redis 服务<br/>redis.service.ts"]
      MQ["内存队列<br/>memory-queue.ts"]
      WS["WebSocket 服务<br/>websocket.service.ts"]
      end
      subgraph "业务层"
      BQP["书籍生成队列处理器<br/>book-queue.processor.ts"]
      FT["容错层<br/>fault-tolerance.ts"]
      BGS["批量生成编排器<br/>book-generator.service.ts"]
      end
      subgraph "前端"
      WSC["WebSocket 工具<br/>websocket.ts"]
      end
      QS --> RS
      QS --> MQ
      BQP --> QS
      BQP --> MQ
      BGS --> WS
      FT --> QS
      WSC --> WS
      

图表来源

  • queue.service.ts:1-347
  • redis.service.ts:1-274
  • memory-queue.ts:1-119
  • book-queue.processor.ts:1-124
  • fault-tolerance.ts:1-387
  • book-generator.service.ts:1-549
  • websocket.service.ts:1-136
  • websocket.ts:1-211

章节来源

  • queue.service.ts:1-347
  • redis.service.ts:1-274
  • memory-queue.ts:1-119
  • book-queue.processor.ts:1-124
  • fault-tolerance.ts:1-387
  • book-generator.service.ts:1-549
  • websocket.service.ts:1-136
  • websocket.ts:1-211

核心组件

  • 队列服务(QueueService)
    • 职责:统一创建/获取队列、添加任务、查询状态、更新进度、统计与控制队列启停
    • 特性:默认不负责重试与超时,交由业务层与容错层处理;支持 Redis 与内存队列双通道
  • 内存队列(MemoryQueue)
    • 职责:在 Redis 不可用时提供本地队列能力,支持并发控制与简单事件
  • 书籍生成队列处理器(book-queue.processor)
    • 职责:绑定 Redis 队列并发与事件,处理书籍生成任务,恢复中断任务
  • 容错层(fault-tolerance)
    • 职责:AI 调用重试(指数退避)、节点超时控制、进度监控、自动恢复
  • 批量生成编排器(BatchGenerationOrchestrator)
    • 职责:按步骤推进内容/音频/视频生成,推送进度,支持取消与超时
  • WebSocket 服务与前端工具
    • 职责:服务端广播进度,前端订阅并渲染

章节来源

  • queue.service.ts:48-347
  • memory-queue.ts:17-119
  • book-queue.processor.ts:16-124
  • fault-tolerance.ts:68-387
  • book-generator.service.ts:45-549
  • websocket.service.ts:1-136
  • websocket.ts:1-211

架构总览

批量处理整体采用“队列解耦 + 业务编排 + 容错保障”的三层架构:

  • 队列层:Bull(Redis)为主,内存队列为回退
  • 业务层:队列处理器负责并发与事件;编排器负责步骤推进与进度推送
  • 容错层:对 AI 调用、节点执行、进度异常进行重试、超时与恢复
  • 前端:通过 WebSocket 实时接收进度与事件

    sequenceDiagram
    participant Client as "客户端"
    participant WS as "WebSocket 服务"
    participant Q as "队列服务"
    participant R as "Redis/Bull"
    participant M as "内存队列"
    participant Proc as "书籍生成处理器"
    participant Biz as "批量编排器"
    Client->>Q : "提交批量任务"
    Q->>R : "尝试添加任务(持久化)"
    alt "Redis 可用"
    R-->>Q : "返回任务ID"
    else "Redis 不可用"
    Q->>M : "回退到内存队列"
    M-->>Q : "返回任务ID"
    end
    Q-->>Client : "返回任务ID"
    Note over Proc,R : "Redis 队列处理器启动并发执行"
    Proc->>Biz : "执行生成步骤"
    Biz-->>WS : "推送进度事件"
    WS-->>Client : "实时进度更新"
    

图表来源

  • queue.service.ts:131-160
  • book-queue.processor.ts:48-82
  • book-generator.service.ts:59-142
  • websocket.service.ts:91-95

详细组件分析

队列系统与任务调度

  • 队列类型与状态
    • 类型:音频生成、视频生成、书籍生成、邮件发送
    • 状态:等待、执行中、完成、失败、延时
  • 任务添加与回退
    • 优先使用 Redis/Bull 队列;失败则回退至内存队列
    • 任务超时在不同队列类型上设置差异化阈值
  • 并发控制与事件
    • Redis 队列处理器设置并发数;内存队列内置并发上限
    • 事件包括完成、失败、进度更新
  • 进度与回调

    • 支持更新任务进度并触发回调
    • 提供进度回调注册与清理

      classDiagram
      class QueueService {
      +isQueueAvailable() bool
      +addTask(type, data, options) Promise<string|null>
      +addAudioGenerationTask(data) Promise<string|null>
      +addVideoGenerationTask(data) Promise<string|null>
      +addBookGenerationTask(data) Promise<string|null>
      +getTaskStatus(type, jobId) Promise<Status>
      +updateProgress(type, jobId, progress, data) Promise<void>
      +onProgress(jobId, cb) void
      +clearProgressCallback(jobId) void
      +getQueueStats(type) Promise<counts>
      +pauseQueue(type) Promise<void>
      +resumeQueue(type) Promise<void>
      +clearQueue(type) Promise<void>
      +closeAll() Promise<void>
      }
      class MemoryQueue {
      +add(name, data) Promise<string>
      +process(concurrency, handler) void
      +hasPendingJobs() bool
      +close() Promise<void>
      }
      QueueService --> MemoryQueue : "回退"
      

图表来源

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

章节来源

  • queue.service.ts:22-347
  • memory-queue.ts:17-119

书籍生成队列处理器

  • 并发与事件
    • Redis 可用时设置并发数,绑定完成/失败/进度事件
    • Redis 不可用时使用内存队列处理器
  • 任务恢复

    • 启动时扫描数据库中处于生成中的书籍,重新入队

      sequenceDiagram
      participant Boot as "应用启动"
      participant Store as "数据库"
      participant Q as "队列服务"
      participant Proc as "书籍生成处理器"
      Boot->>Store : "查询生成中书籍"
      Store-->>Boot : "返回中断任务列表"
      loop "逐个恢复"
      Boot->>Q : "addBookGenerationTask(任务数据)"
      Q-->>Boot : "返回任务ID"
      end
      Proc->>Proc : "启动并发处理器"
      

图表来源

  • book-queue.processor.ts:88-124

章节来源

  • book-queue.processor.ts:48-124

容错层(重试、超时、进度监控、自动恢复)

  • AI 调用重试
    • 最大重试次数、初始/最大延迟、指数退避
    • 失败记录与用户通知
  • 节点超时控制
    • 针对不同节点设定超时阈值
    • 超时后通知用户并进入恢复流程
  • 进度监控
    • 定期检查书籍进度,长时间无进展发出警告,必要时自动恢复
  • 自动恢复

    • 限定最大恢复次数与间隔,失败后标记最终失败

      flowchart TD
      Start(["开始节点执行"]) --> TimeoutCfg["读取节点超时配置"]
      TimeoutCfg --> Exec["执行节点函数"]
      Exec --> Timeout{"是否超时?"}
      Timeout --> |是| NotifyTimeout["通知用户超时"]
      NotifyTimeout --> AutoRecovery["尝试自动恢复"]
      AutoRecovery --> Attempts{"恢复次数 < 最大次数?"}
      Attempts --> |是| Requeue["重新入队"] --> End
      Attempts --> |否| MarkFailed["标记最终失败"] --> End
      Timeout --> |否| Success["执行成功"] --> End(["结束"])
      

图表来源

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

章节来源

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

批量生成编排器(步骤推进与进度推送)

  • 步骤类型:内容生成、音频生成、音频合并、视频生成、视频合并
  • 取消机制:支持设置/检查/清理取消标志
  • 进度推送:通过 WebSocket 广播每一步进度
  • 超时与失败:各步骤设置最大等待时间与轮询间隔,失败时返回失败步骤与错误信息

    sequenceDiagram
    participant Orchestrator as "编排器"
    participant Steps as "步骤集合"
    participant WS as "WebSocket 服务"
    Orchestrator->>Steps : "遍历执行每个步骤"
    loop "每个步骤"
    Orchestrator->>WS : "推送步骤开始进度"
    Orchestrator->>Steps : "执行步骤"
    Steps-->>Orchestrator : "完成/失败"
    alt "完成"
    Orchestrator->>WS : "推送步骤完成进度"
    else "失败"
    Orchestrator-->>Orchestrator : "记录失败步骤与错误"
    Orchestrator->>WS : "推送最终结果"
    end
    end
    Orchestrator->>WS : "推送总体完成"
    

图表来源

  • book-generator.service.ts:77-143
  • websocket.service.ts:91-95

章节来源

  • book-generator.service.ts:45-549
  • websocket.service.ts:91-95

前端进度订阅与渲染

  • WebSocket 工具封装:连接、订阅、取消订阅、重连与重新订阅
  • 前端页面在卸载时取消订阅,避免内存泄漏
  • 通过事件名区分进度与完成事件,按任务ID过滤

章节来源

  • websocket.ts:25-196
  • .plans/audiobook-v2-ux/frontend-dev/task-opt02/progress.md:16-32

依赖关系分析

  • 队列服务依赖 Redis 服务进行可用性判断;当 Redis 不可用时自动回退到内存队列
  • 书籍生成处理器同时兼容 Redis 队列与内存队列
  • 容错层与编排器共同保障任务的稳定性与可观测性
  • WebSocket 服务与前端工具形成闭环,实现进度与事件的实时推送

    graph LR
    Redis["Redis 服务"] --> QS["队列服务"]
    QS --> Bull["Bull/Redis 队列"]
    QS --> MQ["内存队列"]
    BQP["书籍生成处理器"] --> QS
    BQP --> MQ
    FT["容错层"] --> QS
    BGS["批量编排器"] --> WS["WebSocket 服务"]
    WSC["前端 WebSocket 工具"] --> WS
    

图表来源

  • redis.service.ts:43-45
  • queue.service.ts:72-122
  • book-queue.processor.ts:51-82
  • fault-tolerance.ts:13-14
  • book-generator.service.ts:7-11
  • websocket.service.ts:102-133
  • websocket.ts:25-82

章节来源

  • redis.service.ts:43-45
  • queue.service.ts:72-122
  • book-queue.processor.ts:51-82
  • fault-tolerance.ts:13-14
  • book-generator.service.ts:7-11
  • websocket.service.ts:102-133
  • websocket.ts:25-82

性能考量

  • 队列并发与资源限制
    • Redis 队列处理器并发数在书籍生成处理器中固定为 3
    • 内存队列并发上限在内存队列类中固定为 3
  • 超时与重试
    • 不同队列类型设置差异化超时阈值,避免长尾任务占用资源
    • 容错层的指数退避减少对上游服务的压力
  • 性能监控中间件
    • 统计总请求数、平均响应时间、慢请求数量、错误率与端点级指标
    • 可通过路由导出当前指标,便于运维观察

章节来源

  • book-queue.processor.ts:56-76
  • memory-queue.ts:22-22
  • queue.service.ts:166-190
  • fault-tolerance.ts:18-24
  • performance.ts:29-109

故障排查指南

  • 队列不可用
    • 现象:Redis 连接失败或不可用,队列服务标记不可用
    • 处理:检查 Redis 服务状态;确认环境变量;查看队列服务错误日志
  • 任务长时间无响应
    • 现象:进度长时间未更新
    • 处理:容错层会发出警告并在多次无响应后尝试自动恢复;必要时人工干预
  • 任务失败
    • 现象:队列事件失败或编排器步骤失败
    • 处理:查看失败原因与错误日志;根据容错层通知进行重试或恢复
  • 前端进度不更新
    • 现象:WebSocket 断开或未收到事件
    • 处理:检查前端 WebSocket 工具连接状态与订阅事件;确认服务端 WebSocket 服务运行

章节来源

  • queue.service.ts:54-58
  • redis.service.ts:24-38
  • fault-tolerance.ts:188-261
  • book-queue.processor.ts:66-74
  • websocket.ts:163-178

结论

本批量处理机制通过“队列解耦 + 业务编排 + 容错保障”的架构实现了高可用、可观测与可扩展的批处理能力。Redis/Bull 作为主队列,内存队列提供强健的回退;容错层覆盖重试、超时与自动恢复;编排器与 WebSocket 形成清晰的进度闭环。配合性能监控中间件与合理的并发/超时策略,系统在复杂业务场景下具备良好的稳定性与可维护性。

附录

配置参数说明

  • Redis 连接与可用性
    • REDIS_HOST、REDIS_PORT、REDIS_PASSWORD、REDIS_DB
    • Redis 服务提供连接状态检测与错误处理
  • 队列默认行为
    • 默认移除完成/失败记录数量(用于统计)
    • 队列空闲检测与停滞任务处理
  • 模型与服务配置
    • JWT、DashScope TTS、上传大小限制等
  • WebSocket
    • 服务端升级路径 /ws,客户端通过 wsManager 订阅事件

章节来源

  • redis.service.ts:8-22
  • queue.service.ts:87-95
  • index.ts:70-117
  • websocket.service.ts:102-133
  • websocket.ts:31-32

扩展接口设计

  • 队列服务扩展点
    • 新增队列类型:在枚举与添加方法中扩展
    • 新增任务选项:在任务添加时传入自定义选项
  • 内存队列扩展点
    • 调整并发上限、自定义处理器
  • 容错层扩展点
    • 新增节点超时阈值、新增重试策略、新增自动恢复策略
  • WebSocket 扩展点
    • 新增事件类型与前端订阅

章节来源

  • queue.service.ts:22-347
  • memory-queue.ts:48-52
  • fault-tolerance.ts:17-51
  • websocket.service.ts:67-95

数据一致性与生命周期管理

  • 一致性
    • 任务状态通过队列事件与数据库状态同步
    • 容错层记录失败与恢复历史,便于审计
  • 生命周期
    • 任务创建、排队、执行、完成/失败、清理
    • 中断任务恢复:启动时扫描数据库并重新入队
  • 超时与重试
    • 队列级别超时与业务级别超时结合
    • 容错层指数退避与最大重试次数控制

章节来源

  • book-queue.processor.ts:88-124
  • fault-tolerance.ts:18-24
  • queue.service.ts:166-190