2026-06-03 10:30:42 +08:00
|
|
|
---
|
|
|
|
|
tags: [consumer-producer, bridge-pattern, decoupling, go, message-queue]
|
|
|
|
|
create time: 2026-06-03 10:50
|
|
|
|
|
---
|
|
|
|
|
|
2026-06-03 10:12:49 +08:00
|
|
|
# 11. Consumer-Producer 桥接模式
|
|
|
|
|
|
2026-06-03 10:30:42 +08:00
|
|
|
## 概述
|
2026-06-03 10:12:49 +08:00
|
|
|
|
2026-06-03 10:30:42 +08:00
|
|
|
Consumer 结构体桥接 TaskQueue 和 WorkerPool,实现生产者与消费者的彻底解耦。
|
|
|
|
|
|
|
|
|
|
## 正文
|
2026-06-03 10:12:49 +08:00
|
|
|
|
2026-06-03 10:30:42 +08:00
|
|
|
### 架构总览
|
2026-06-03 10:12:49 +08:00
|
|
|
|
|
|
|
|
```mermaid
|
|
|
|
|
flowchart LR
|
2026-06-03 10:30:42 +08:00
|
|
|
subgraph Producer["Producer"]
|
|
|
|
|
A["Generate Handler"]
|
2026-06-03 10:12:49 +08:00
|
|
|
end
|
|
|
|
|
|
|
|
|
|
subgraph Queue["TaskQueue 接口"]
|
2026-06-03 10:30:42 +08:00
|
|
|
B[("Memory Queue")]
|
|
|
|
|
C[("RabbitMQ Queue")]
|
2026-06-03 10:12:49 +08:00
|
|
|
end
|
|
|
|
|
|
|
|
|
|
subgraph Bridge["Consumer 桥接层"]
|
2026-06-03 10:30:42 +08:00
|
|
|
D{{"Consumer"}}
|
2026-06-03 10:12:49 +08:00
|
|
|
end
|
|
|
|
|
|
|
|
|
|
subgraph Pool["WorkerPool"]
|
2026-06-03 10:30:42 +08:00
|
|
|
E["Worker 1"]
|
|
|
|
|
F["Worker 2"]
|
|
|
|
|
G["Worker N"]
|
2026-06-03 10:12:49 +08:00
|
|
|
end
|
|
|
|
|
|
|
|
|
|
subgraph Pipeline["业务逻辑"]
|
2026-06-03 10:30:42 +08:00
|
|
|
H["Eino Pipeline"]
|
2026-06-03 10:12:49 +08:00
|
|
|
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
|
|
|
|
|
```
|
|
|
|
|
|
2026-06-03 10:30:42 +08:00
|
|
|
### 核心结构体
|
2026-06-03 10:12:49 +08:00
|
|
|
|
|
|
|
|
`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 层注入实际业务逻辑 |
|
|
|
|
|
|
|
|
|
|
---
|
|
|
|
|
|
2026-06-03 10:30:42 +08:00
|
|
|
### 工作流程
|
2026-06-03 10:12:49 +08:00
|
|
|
|
2026-06-03 10:30:42 +08:00
|
|
|
#### 启动消费循环
|
2026-06-03 10:12:49 +08:00
|
|
|
|
|
|
|
|
```mermaid
|
|
|
|
|
sequenceDiagram
|
|
|
|
|
participant Main as main.go
|
2026-06-03 10:30:42 +08:00
|
|
|
participant Consumer as Consumer
|
2026-06-03 10:12:49 +08:00
|
|
|
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 管线
|
|
|
|
|
|
|
|
|
|
---
|
|
|
|
|
|
2026-06-03 10:30:42 +08:00
|
|
|
### 解耦的三层设计
|
2026-06-03 10:12:49 +08:00
|
|
|
|
|
|
|
|
```mermaid
|
|
|
|
|
flowchart TB
|
2026-06-03 10:30:42 +08:00
|
|
|
subgraph "消息源"
|
2026-06-03 10:12:49 +08:00
|
|
|
Q["TaskQueue 接口\n(Memory / RabbitMQ)"]
|
|
|
|
|
end
|
2026-06-03 10:30:42 +08:00
|
|
|
subgraph "桥接"
|
2026-06-03 10:12:49 +08:00
|
|
|
C["Consumer\n(只关心 消费→提交)"]
|
|
|
|
|
end
|
2026-06-03 10:30:42 +08:00
|
|
|
subgraph "执行引擎"
|
2026-06-03 10:12:49 +08:00
|
|
|
P["WorkerPool\n(只关心 并发控制)"]
|
|
|
|
|
end
|
2026-06-03 10:30:42 +08:00
|
|
|
subgraph "业务逻辑"
|
2026-06-03 10:12:49 +08:00
|
|
|
H["TaskHandler 回调\n(RunFromTaskMessage)"]
|
|
|
|
|
end
|
|
|
|
|
|
|
|
|
|
Q --> C --> P --> H
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
| 组件 | 知道什么 | 不知道什么 |
|
|
|
|
|
|------|----------|------------|
|
|
|
|
|
| **TaskQueue** | 消息的存储与投递 | WorkerPool 的存在 |
|
|
|
|
|
| **WorkerPool** | 任务的并发执行 | 消息来自哪个队列 |
|
|
|
|
|
| **Consumer** | 如何桥接两者 | 具体的业务逻辑 |
|
|
|
|
|
| **TaskHandler** | 生成管线的执行 | 消息来自内存还是 RabbitMQ |
|
|
|
|
|
|
|
|
|
|
> **设计哲学**:每个组件只关心自己的职责边界,可独立替换、测试、扩展。
|
|
|
|
|
|
|
|
|
|
---
|
|
|
|
|
|
2026-06-03 10:30:42 +08:00
|
|
|
### Handler 层注入
|
2026-06-03 10:12:49 +08:00
|
|
|
|
|
|
|
|
`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 管线。
|
|
|
|
|
|
|
|
|
|
---
|
|
|
|
|
|
2026-06-03 10:30:42 +08:00
|
|
|
### 信号处理与优雅关闭
|
2026-06-03 10:12:49 +08:00
|
|
|
|
|
|
|
|
```mermaid
|
|
|
|
|
sequenceDiagram
|
2026-06-03 10:30:42 +08:00
|
|
|
participant OS as OS
|
2026-06-03 10:12:49 +08:00
|
|
|
participant Main as main.go
|
2026-06-03 10:30:42 +08:00
|
|
|
participant Consumer as Consumer
|
2026-06-03 10:12:49 +08:00
|
|
|
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 / 内存资源
|
|
|
|
|
|
|
|
|
|
## 关联文档
|
|
|
|
|
|
2026-06-03 10:30:42 +08:00
|
|
|
- [[03-任务队列]] — Consumer 的消息来源
|
|
|
|
|
- [[02-协程池]] — Consumer 的执行引擎
|
|
|
|
|
- [[12-三级降级策略]] — Consumer 不参与降级,降级在 Handler 层
|
|
|
|
|
- [[00-索引]] — 返回文档总览
|