Files
cs-note/hzh/Gen2D/11-consumer-producer.md
T

186 lines
4.9 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 11. Consumer-Producer 桥接模式
> **一句话概括**:`Consumer` 结构体桥接 `TaskQueue` 和 `WorkerPool`,实现生产者与消费者的彻底解耦。
---
## 架构总览
```mermaid
flowchart LR
subgraph Producer["生产者"]
A[Generate Handler]
end
subgraph Queue["TaskQueue 接口"]
B((Memory\nQueue))
C((RabbitMQ\nQueue))
end
subgraph Bridge["Consumer 桥接层"]
D{{"Consumer\n(bridge)"}}
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
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 "第 1 层:消息源"
Q["TaskQueue 接口\n(Memory / RabbitMQ)"]
end
subgraph "第 2 层:桥接"
C["Consumer\n(只关心 消费→提交)"]
end
subgraph "第 3 层:执行引擎"
P["WorkerPool\n(只关心 并发控制)"]
end
subgraph "第 4 层:业务逻辑"
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 操作系统
participant Main as main.go
participant 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 - 任务队列](03-task-queue.md) | Consumer 的消息来源 |
| [02 - 协程池](02-worker-pool.md) | Consumer 的执行引擎 |
| [12 - 三级降级策略](12-three-tier-fallback.md) | Consumer 不参与降级,降级在 Handler 层 |
| [00 - 索引](00-index.md) | 返回文档总览 |