数据流设计.md 21 KB

数据流设计

本文档引用的文件

  • server/src/app.ts
  • server/src/modules/book-generator/index.ts
  • server/src/modules/book-generator/book-generator.controller.ts
  • server/src/modules/book-generator/book-generator.service.ts
  • server/src/modules/book-generator/book-queue.processor.ts
  • server/src/modules/tts/tts.controller.ts
  • server/src/modules/tts/tts.service.ts
  • server/src/modules/video-generator/video-generator.controller.ts
  • server/src/modules/video-generator/video-generator.service.ts
  • server/src/services/queue.service.ts
  • server/src/services/websocket.service.ts
  • server/src/services/storage.service.ts
  • server/src/services/oss.service.ts
  • server/src/models/index.ts

目录

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

简介

本文件面向AI有声书生成平台,系统性梳理从用户输入到最终输出的完整数据流,重点覆盖:

  • 异步数据处理流程:任务队列调度机制、LangGraph工作流的数据传递、TTS服务的音频生成流程
  • 数据在不同系统组件间的传输方式:HTTP请求的数据交换、WebSocket的实时数据推送、文件上传下载的数据传输
  • 数据的持久化策略:数据库写入时机、缓存更新机制、文件存储位置
  • 关键业务场景的数据流图:书籍生成、音频合成、视频制作

项目结构

后端采用Koa应用,集中注册各类模块路由,并在启动时初始化数据库、Redis、存储与WebSocket服务,随后启动书籍生成队列处理器与恢复中断任务。

graph TB
A["应用入口<br/>server/src/app.ts"] --> B["中间件与日志<br/>错误、性能、安全、限流"]
A --> C["静态资源服务<br/>音频/视频上传目录"]
A --> D["路由注册<br/>认证、TTS、视频、书籍生成等"]
A --> E["服务初始化<br/>数据库、Redis、存储、WebSocket"]
A --> F["队列处理器<br/>书籍生成队列"]

图表来源

  • server/src/app.ts:133-194

章节来源

  • server/src/app.ts:57-131

核心组件

  • 应用与路由:统一创建HTTP服务、注册路由、挂载静态资源、健康检查与指标接口
  • 书籍生成:LangGraph工作流编排、批量生成编排器、队列处理器与中断恢复
  • TTS服务:异步音频生成、提供商选择与降级、分段与合并、存储上传、WebSocket事件推送
  • 视频生成:视频项目管理、FFmpeg生成、进度与状态持久化、WebSocket事件推送
  • 队列服务:Bull队列封装、并发控制、进度回调、Redis/内存双栈
  • 存储服务:OSS与本地存储无缝切换、上传/下载/删除/签名URL
  • WebSocket服务:客户端连接管理、事件广播、音频/视频生成完成通知

章节来源

  • server/src/app.ts:26-55
  • server/src/modules/book-generator/index.ts:60-77
  • server/src/modules/book-generator/book-generator.service.ts:45-143
  • server/src/modules/tts/tts.service.ts:200-280
  • server/src/modules/video-generator/video-generator.service.ts:157-312
  • server/src/services/queue.service.ts:48-346
  • server/src/services/storage.service.ts:13-278
  • server/src/services/websocket.service.ts:6-136

架构总览

系统采用“HTTP API + 异步队列 + WebSocket推送”的架构,数据在模块间通过请求体、响应体、文件系统与数据库进行传递,存储层支持OSS与本地两种模式。

graph TB
subgraph "客户端"
FE["前端/小程序"]
end
subgraph "后端"
HTTP["HTTP服务<br/>server/src/app.ts"]
WS["WebSocket服务<br/>server/src/services/websocket.service.ts"]
QUEUE["队列服务<br/>server/src/services/queue.service.ts"]
STORE["存储服务<br/>server/src/services/storage.service.ts"]
DB["数据库<br/>server/src/models/index.ts"]
end
subgraph "业务模块"
BG["书籍生成<br/>book-generator.*"]
TTS["TTS服务<br/>tts.*"]
VG["视频生成<br/>video-generator.*"]
end
FE --> HTTP
HTTP --> BG
HTTP --> TTS
HTTP --> VG
BG --> QUEUE
BG --> DB
TTS --> DB
TTS --> STORE
VG --> DB
VG --> STORE
BG --> WS
TTS --> WS
VG --> WS

图表来源

  • server/src/app.ts:100-129
  • server/src/services/websocket.service.ts:102-136
  • server/src/services/queue.service.ts:72-122
  • server/src/services/storage.service.ts:43-93
  • server/src/models/index.ts:5-13

详细组件分析

书籍生成数据流(LangGraph工作流)

  • 用户通过API触发书籍生成,系统根据书籍规模与类型推导大纲层级,创建书籍并尝试加入队列
  • 队列处理器异步执行LangGraph生成器,更新书籍状态并在完成后推送进度
  • 批量生成编排器支持多步骤串联:内容生成→音频生成→音频合并→视频生成→视频合并,每步通过WebSocket推送进度

    sequenceDiagram
    participant U as "用户"
    participant API as "书籍生成API<br/>book-generator.controller.ts"
    participant Svc as "编排服务<br/>book-generator.service.ts"
    participant Q as "队列处理器<br/>book-queue.processor.ts"
    participant LG as "LangGraph生成器<br/>book-generator/index.ts"
    participant DB as "数据库<br/>models/index.ts"
    participant WS as "WebSocket<br/>websocket.service.ts"
    U->>API : POST /api/book-generator/langgraph/books
    API->>Svc : 创建书籍并加入队列
    Svc->>Q : addBookGenerationTask
    Q->>LG : 处理任务并执行生成
    LG->>DB : 更新书籍状态/大纲/章节
    LG-->>WS : 推送生成进度/完成事件
    API-->>U : 返回任务状态
    

图表来源

  • server/src/modules/book-generator/book-generator.controller.ts:383-530
  • server/src/modules/book-generator/book-queue.processor.ts:16-83
  • server/src/modules/book-generator/index.ts:60-77
  • server/src/services/websocket.service.ts:91-95
  • server/src/models/index.ts:5-13

章节来源

  • server/src/modules/book-generator/book-generator.controller.ts:24-119
  • server/src/modules/book-generator/book-generator.service.ts:45-143
  • server/src/modules/book-generator/book-queue.processor.ts:48-124
  • server/src/modules/book-generator/index.ts:60-104

批量生成编排器(多步骤流水线)

  • 编排器按步骤推进,每个步骤完成后通过WebSocket推送进度
  • 步骤包括:内容生成(LangGraph)、音频生成(异步队列/轮询)、音频合并(章节聚合)、视频生成(FFmpeg)、视频合并(简化处理)
  • 支持取消标志,中途取消会清理状态并推送失败事件

    flowchart TD
    Start(["开始"]) --> CheckCancel["检查取消标志"]
    CheckCancel --> Step1["内容生成<br/>LangGraph生成大纲/内容"]
    Step1 --> Step2["音频生成<br/>异步生成+轮询完成"]
    Step2 --> MergeA["音频合并<br/>按父章节合并"]
    MergeA --> Step3["视频生成<br/>FFmpeg生成视频"]
    Step3 --> MergeV["视频合并<br/>简化处理"]
    MergeV --> Done(["完成"])
    Step1 -.->|失败| Fail["失败并清理"]
    Step2 -.->|失败| Fail
    MergeA -.->|失败| Fail
    Step3 -.->|失败| Fail
    MergeV -.->|失败| Fail
    

图表来源

  • server/src/modules/book-generator/book-generator.service.ts:77-143
  • server/src/modules/book-generator/book-generator.service.ts:149-217
  • server/src/modules/book-generator/book-generator.service.ts:222-285
  • server/src/modules/book-generator/book-generator.service.ts:289-359
  • server/src/modules/book-generator/book-generator.service.ts:364-453
  • server/src/modules/book-generator/book-generator.service.ts:458-528

章节来源

  • server/src/modules/book-generator/book-generator.service.ts:45-143

TTS服务数据流(异步音频生成)

  • 用户提交文本、音色、参数,服务创建AudioRecord记录并异步生成
  • 文本分段(阿里云限制)或直传(MiniMax),并发生成后合并,上传至存储(OSS或本地)
  • 生成完成后更新数据库状态并推送WebSocket事件,支持预览与批量下载

    sequenceDiagram
    participant U as "用户"
    participant API as "TTS API<br/>tts.controller.ts"
    participant Svc as "TTS服务<br/>tts.service.ts"
    participant Prov as "TTS提供商<br/>aliyun/minimax/mock"
    participant Store as "存储服务<br/>storage.service.ts"
    participant DB as "数据库<br/>models/index.ts"
    participant WS as "WebSocket<br/>websocket.service.ts"
    U->>API : POST /api/tts/generate
    API->>Svc : generateAudio(userId,text,voice,...)
    Svc->>DB : 创建AudioRecord(processing)
    Svc->>Prov : 分段/直传合成
    Prov-->>Svc : 音频片段/云端URL
    Svc->>Store : 合并后上传(oss/local)
    Store-->>Svc : 返回文件URL
    Svc->>DB : 更新AudioRecord(completed)
    Svc->>WS : 推送音频生成完成事件
    API-->>U : 返回任务已创建
    

图表来源

  • server/src/modules/tts/tts.controller.ts:53-127
  • server/src/modules/tts/tts.service.ts:200-280
  • server/src/modules/tts/tts.service.ts:285-542
  • server/src/services/storage.service.ts:43-93
  • server/src/services/websocket.service.ts:70-76
  • server/src/models/index.ts:5-13

章节来源

  • server/src/modules/tts/tts.controller.ts:13-127
  • server/src/modules/tts/tts.service.ts:200-542

视频生成数据流(FFmpeg流水线)

  • 用户创建视频项目或从书籍章节生成项目,系统校验素材并生成视频
  • FFmpeg生成完成后更新项目状态,若关联章节则更新章节视频URL并推送完成事件
  • 支持带/不带背景音乐两种生成路径

    sequenceDiagram
    participant U as "用户"
    participant API as "视频API<br/>video-generator.controller.ts"
    participant Svc as "视频服务<br/>video-generator.service.ts"
    participant FF as "FFmpeg生成器"
    participant Store as "存储服务<br/>storage.service.ts"
    participant DB as "数据库<br/>models/index.ts"
    participant WS as "WebSocket<br/>websocket.service.ts"
    U->>API : POST /api/video/projects/ : id/generate
    API->>Svc : generateVideoForProject(projectId)
    Svc->>DB : 更新状态为processing
    Svc->>FF : 生成视频(带/不带BGM)
    FF-->>Svc : 输出文件路径/时长/大小
    Svc->>DB : 更新状态为completed并写入URL
    Svc->>WS : 推送视频生成完成事件
    API-->>U : 返回输出URL/时长/大小
    

图表来源

  • server/src/modules/video-generator/video-generator.controller.ts:111-129
  • server/src/modules/video-generator/video-generator.service.ts:157-312
  • server/src/services/websocket.service.ts:81-86
  • server/src/models/index.ts:5-13

章节来源

  • server/src/modules/video-generator/video-generator.controller.ts:24-242
  • server/src/modules/video-generator/video-generator.service.ts:157-312

队列服务与并发控制

  • 队列服务封装Bull,支持Redis队列与内存队列双栈,自动降级
  • 书籍生成队列最大并发3,处理完成后更新书籍状态
  • 支持任务状态查询、进度回调、统计与暂停/恢复

    classDiagram
    class QueueService {
    +addTask(queueType,data,options) Promise~string|null~
    +addBookGenerationTask(data) Promise~string|null~
    +getTaskStatus(queueType,jobId) Promise
    +updateProgress(queueType,jobId,progress,data) Promise
    +onProgress(jobId,callback) void
    +getQueueStats(queueType) Promise
    }
    class BookQueueProcessor {
    +initBookGenerationQueue() void
    +resumeInterruptedTasks() void
    }
    QueueService <.. BookQueueProcessor : "使用"
    

图表来源

  • server/src/services/queue.service.ts:48-346
  • server/src/modules/book-generator/book-queue.processor.ts:48-124

章节来源

  • server/src/services/queue.service.ts:48-346
  • server/src/modules/book-generator/book-queue.processor.ts:48-124

存储与文件传输

  • 存储服务支持OSS与本地存储无缝切换,提供上传/下载/删除/签名URL能力
  • TTS与视频生成均通过存储服务统一分发,OSS模式下支持CDN加速
  • 文件上传通过koa-body中间件接收multipart/form-data

    graph LR
    Svc["业务服务<br/>tts.service.ts / video-generator.service.ts"] --> Store["存储服务<br/>storage.service.ts"]
    Store --> OSS["OSS服务<br/>oss.service.ts"]
    Store --> Local["本地文件系统"]
    OSS --> CDN["CDN加速(可选)"]
    

图表来源

  • server/src/services/storage.service.ts:43-93
  • server/src/services/oss.service.ts:38-117
  • server/src/modules/tts/tts.service.ts:406-436
  • server/src/modules/video-generator/video-generator.service.ts:257-268

章节来源

  • server/src/services/storage.service.ts:13-278
  • server/src/services/oss.service.ts:13-256
  • server/src/app.ts:76-83

WebSocket实时推送

  • 服务启动时初始化WebSocket,只接受/ws路径连接
  • 推送音频/视频生成完成事件与批量生成进度事件
  • 前端可通过clientId区分接收范围

    sequenceDiagram
    participant FE as "前端客户端"
    participant WS as "WebSocket服务<br/>websocket.service.ts"
    participant Svc as "业务服务"
    FE->>WS : 连接 /ws?clientId=...
    WS-->>FE : connected
    Svc->>WS : pushAudioGenerationComplete(...)
    WS-->>FE : 广播事件
    Svc->>WS : pushBatchGenerationProgress(...)
    WS-->>FE : 广播事件
    

图表来源

  • server/src/services/websocket.service.ts:102-136
  • server/src/services/websocket.service.ts:70-95

章节来源

  • server/src/services/websocket.service.ts:6-136

依赖关系分析

  • 控制器层依赖服务层,服务层依赖数据库与存储服务
  • 队列服务独立于业务服务,通过任务数据解耦
  • WebSocket服务与业务服务松耦合,通过事件广播解耦

    graph TB
    Ctrl["控制器层<br/>*.controller.ts"] --> Svc["服务层<br/>*.service.ts"]
    Svc --> DB["数据库<br/>models/index.ts"]
    Svc --> Store["存储服务<br/>storage.service.ts"]
    Svc --> Queue["队列服务<br/>queue.service.ts"]
    Svc --> WS["WebSocket服务<br/>websocket.service.ts"]
    Queue --> LG["LangGraph生成器<br/>book-generator/index.ts"]
    

图表来源

  • server/src/modules/book-generator/book-generator.controller.ts:18-199
  • server/src/modules/tts/tts.controller.ts:10-274
  • server/src/modules/video-generator/video-generator.controller.ts:20-244
  • server/src/services/queue.service.ts:48-346
  • server/src/services/websocket.service.ts:6-136
  • server/src/services/storage.service.ts:13-278
  • server/src/models/index.ts:5-13

章节来源

  • server/src/app.ts:100-129

性能考量

  • 队列并发:书籍生成队列最大并发3,避免资源争用
  • 文本分段:TTS服务对阿里云进行分段,减少单次请求压力
  • 并发合成:MiniMax异步轮询较长时降低并发至1,其他提供商为2
  • 存储上传:统一通过存储服务上传,支持OSS直传与CDN加速
  • WebSocket:事件广播轻量,前端按bookId过滤,避免过多无效推送

故障排查指南

  • 队列不可用:Redis连接失败时自动降级为内存队列,检查Redis配置与连通性
  • 任务超时:僵尸任务检测(空目录超过2分钟)会更新数据库状态为失败
  • 存储失败:OSS连接失败时可切换为本地存储,检查OSS配置与凭证
  • WebSocket连接:确认只接受/ws路径,客户端需携带clientId参数
  • TTS配额:检查用户音频分钟配额,额度不足会阻止生成

章节来源

  • server/src/services/queue.service.ts:72-122
  • server/src/modules/tts/tts.service.ts:574-597
  • server/src/services/storage.service.ts:252-272
  • server/src/services/websocket.service.ts:104-130

结论

本平台通过HTTP API、异步队列与WebSocket实现了高并发、可扩展的AI有声书生成体系。数据在模块间以清晰的职责边界传递,存储与数据库持久化策略明确,具备良好的可观测性与可维护性。建议在生产环境中启用Redis队列、OSS存储与CDN加速,并结合WebSocket事件实现前端实时反馈。

附录

  • 关键API路径
    • 书籍生成:/api/book-generator/langgraph/books
    • 批量生成:/api/book-generator/books/:id/batch-generate
    • TTS生成:/api/tts/generate
    • 视频生成:/api/video/projects/:id/generate
  • 关键事件
    • audio_generation_complete
    • video_generation_complete
    • batch_generation_progress