--- tags: [consumer-producer, bridge-pattern, decoupling, go, message-queue] create time: 2026-06-03 10:50 --- # 11. Consumer-Producer 桥接模式 ## 概述 Consumer 结构体桥接 TaskQueue 和 WorkerPool,实现生产者与消费者的彻底解耦。 ## 正文 ### 架构总览 ```mermaid flowchart LR subgraph Producer["Producer"] A["Generate Handler"] end subgraph Queue["TaskQueue 接口"] B[("Memory Queue")] C[("RabbitMQ Queue")] end subgraph Bridge["Consumer 桥接层"] D{{"Consumer"}} end subgraph Pool["WorkerPool"] E["Worker 1"] F["Worker 2"] G["Worker N"] end subgraph Pipeline["业务逻辑"] H["Eino Pipeline"] end A -->|"Submit(msg)"| B A -->|"Submit(msg)"| C B -->|"Consume()"| D C -->|"Consume()"| D D -->|"Submit(task)"| E D -->|"Submit(task)"| F D -->|"Submit(task)"| G E --> H F --> H G --> H style D fill:#f9a825,stroke:#333,color:#000 ``` ### 核心结构体 `Consumer` 是整个任务调度体系的**桥梁**,它只做一件事:从队列取消息,提交到协程池。 ```go // worker/consumer.go type Consumer struct { queue taskqueue.TaskQueue // 可插拔队列接口 pool *workerpool.Pool // 有界协程池 handler TaskHandler // 业务逻辑注入点 logger *slog.Logger } ``` | 字段 | 类型 | 职责 | |------|------|------| | `queue` | `TaskQueue` 接口 | 消息来源,支持 Memory / RabbitMQ 替换 | | `pool` | `*workerpool.Pool` | 并发执行引擎,控制单机并行度 | | `handler` | `TaskHandler` | 回调函数,由 handler 层注入实际业务逻辑 | --- ### 工作流程 #### 启动消费循环 ```mermaid sequenceDiagram participant Main as main.go participant Consumer as Consumer participant Queue as TaskQueue participant Pool as WorkerPool participant Handler as RunFromTaskMessage Main->>Consumer: Start(ctx) loop 持续消费 Consumer->>Queue: Consume(ctx, callback) Queue-->>Consumer: TaskMessage Consumer->>Consumer: processMessage() Consumer->>Pool: Submit(workerpool.Task) Pool-->>Consumer: accepted Pool->>Handler: Fn(ctx) Handler->>Handler: runPipelineBg() end ``` 四步循环: 1. **消费** — `Consumer` 调用 `queue.Consume(ctx, handler)`,阻塞等待消息 2. **转换** — 将 `taskqueue.TaskMessage` 包装为 `workerpool.Task` 3. **提交** — 调用 `pool.Submit(task)` 送入协程池执行 4. **执行** — `Task.Fn` 回调实际的 `RunFromTaskMessage`,驱动 Eino 管线 --- ### 解耦的三层设计 ```mermaid flowchart TB subgraph "消息源" Q["TaskQueue 接口\n(Memory / RabbitMQ)"] end subgraph "桥接" C["Consumer\n(只关心 消费→提交)"] end subgraph "执行引擎" P["WorkerPool\n(只关心 并发控制)"] end subgraph "业务逻辑" H["TaskHandler 回调\n(RunFromTaskMessage)"] end Q --> C --> P --> H ``` | 组件 | 知道什么 | 不知道什么 | |------|----------|------------| | **TaskQueue** | 消息的存储与投递 | WorkerPool 的存在 | | **WorkerPool** | 任务的并发执行 | 消息来自哪个队列 | | **Consumer** | 如何桥接两者 | 具体的业务逻辑 | | **TaskHandler** | 生成管线的执行 | 消息来自内存还是 RabbitMQ | > **设计哲学**:每个组件只关心自己的职责边界,可独立替换、测试、扩展。 --- ### Handler 层注入 `Consumer` 不硬编码业务逻辑,而是通过 `TaskHandler` 函数签名由外部注入: ```go // worker/consumer.go — 定义 type TaskHandler func(ctx context.Context, msg taskqueue.TaskMessage) // main.go — 注入 consumer := worker.NewConsumer(tq, pool, handler.RunFromTaskMessage) ``` `RunFromTaskMessage` 负责将队列消息还原为 `GenerateRequest`,再调用 `runPipelineBg` 驱动完整的 Eino 管线。 --- ### 信号处理与优雅关闭 ```mermaid sequenceDiagram participant OS as OS participant Main as main.go participant Consumer as Consumer participant Pool as WorkerPool OS->>Main: SIGINT / SIGTERM Main->>Consumer: cancel() 停止消费 Note over Consumer: 不再接收新消息 Main->>Pool: Shutdown(30s timeout) Note over Pool: 等待正在执行的任务完成 Pool-->>Main: true (正常) / false (超时) Main->>Main: 进程退出 ``` 关闭顺序至关重要: 1. **先停 Consumer** — 不再从队列拉取新消息 2. **再关 WorkerPool** — 等待已提交的任务执行完毕(最多 30 秒) 3. **最后关闭队列连接** — 释放 RabbitMQ / 内存资源 ## 关联文档 - [[03-任务队列]] — Consumer 的消息来源 - [[02-协程池]] — Consumer 的执行引擎 - [[12-三级降级策略]] — Consumer 不参与降级,降级在 Handler 层 - [[00-索引]] — 返回文档总览