Files
cs-note/hzh/Gen2D/03-任务队列.md
T

265 lines
7.6 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.
---
tags: [task-queue, plugin-architecture, memory-queue, rabbitmq, go, interface-pattern]
create time: 2026-06-03 10:10
---
# 03. 任务队列 (Task Queue)
## 概述
可插拔任务队列接口,支持 Memory 和 RabbitMQ 双实现,通过工厂模式一行切换,满足不同部署环境的需求。
> **一句话概括**:可插拔任务队列接口,支持 Memory 和 RabbitMQ 双实现,通过工厂模式一行切换。
## 正文
### 架构设计
```mermaid
graph TB
subgraph Producer["生产者"]
HANDLER["Handler.Generate"]
end
subgraph Interface["TaskQueue接口"]
SUBMIT["Submitctx msg"]
CONSUME["Consumectx handler"]
CLOSE["Close"]
end
subgraph Memory["MemoryQueue"]
MEM_CH["chan TaskMessage<br/>buffered"]
MEM_SUBMIT["阻塞写入"]
MEM_CONSUME["持续消费"]
end
subgraph RabbitMQ["RabbitMQQueue"]
RMQ_PUB["channel.Publish<br/>Persistent"]
RMQ_CONSUME["channel.Consume<br/>Manual ACK"]
end
HANDLER --> SUBMIT
SUBMIT --> Memory
SUBMIT --> RabbitMQ
CONSUME --> Memory
CONSUME --> RabbitMQ
Memory --- MEM_CH
RabbitMQ --- RMQ_PUB
```
### TaskQueue 接口
```go
type TaskQueue interface {
Submit(ctx context.Context, msg TaskMessage) error
Consume(ctx context.Context, handler func(TaskMessage) error) error
Close() error
}
```
| 方法 | 语义 | 错误处理 |
|------|------|----------|
| `Submit` | 提交任务到队列 | 内存队列 buffer 满时阻塞;RabbitMQ 发送失败时返回错误 |
| `Consume` | 持续消费任务,直到 ctx 取消 | handler 返回错误时,内存队列丢弃,RabbitMQ NACK 重试 |
| `Close` | 关闭连接,释放资源 | 返回 `errors.Join` 聚合错误 |
### 工厂模式
通过配置驱动,一行切换队列实现:
```go
tq, err := taskqueue.New(cfg.TaskQueue)
```
```go
func New(cfg config.TaskQueueConfig) (TaskQueue, error) {
switch cfg.Driver {
case "rabbitmq":
return NewRabbitMQQueue(cfg.RabbitMQ)
default:
return NewMemoryQueue(cfg.Memory.BufferSize), nil
}
}
```
| `cfg.Driver` | 实现 | 适用场景 |
|--------------|------|----------|
| `"memory"` (默认) | `MemoryQueue` | 单机开发、演示环境 |
| `"rabbitmq"` | `RabbitMQQueue` | 多机生产部署 |
### TaskMessage 消息结构
```go
type TaskMessage struct {
TaskID string `json:"task_id"`
UserID string `json:"user_id"`
ProjectID string `json:"project_id"`
Input service.PipelineInput `json:"input"`
CreatedAt time.Time `json:"created_at"`
RetryCount int `json:"retry_count"`
}
```
**关键设计**:
- `Input` 直接嵌入 `PipelineInput`,Consumer 反序列化后即可传入管线
- `RetryCount` 供 RabbitMQ 实现判断是否超过最大重试次数
- JSON 序列化,兼容内存队列和 RabbitMQ 两种传输
### MemoryQueue 实现
基于 Go channel 的内存队列,零外部依赖。
**核心逻辑**:
```go
type MemoryQueue struct {
ch chan TaskMessage // 有界缓冲
closed atomic.Bool
closeCh chan struct{}
}
```
#### Submit
```go
func (q *MemoryQueue) Submit(ctx context.Context, msg TaskMessage) error {
select {
case q.ch <- msg: // 成功入队
return nil
case <-ctx.Done(): // 上下文取消
return ctx.Err()
case <-q.closeCh: // 队列已关闭
return ErrQueueClosed
}
}
```
> [!tip] 阻塞语义
>
> MemoryQueue 的 Submit 是阻塞的——当 buffer 满时,调用方会阻塞直到有空位或 ctx 取消。这与 WorkerPool 的非阻塞拒绝形成对比。
#### Consume
```go
func (q *MemoryQueue) Consume(ctx context.Context, handler func(TaskMessage) error) error {
for {
select {
case msg := <-q.ch: // 取出消息
if err := handler(msg); err != nil {
// handler 失败,记录日志并丢弃(不重试)
slog.Error("handler failed, dropping message", ...)
}
case <-ctx.Done(): // 上下文取消
return ctx.Err()
case <-q.closeCh: // 队列已关闭
return ErrQueueClosed
}
}
}
```
**MemoryQueue 特性**:
| 特性 | 行为 |
|------|------|
| 持久化 | 无(进程重启后任务丢失) |
| 重试 | 无(handler 错误直接丢弃) |
| 背压 | 阻塞直到有空位 |
| 依赖 | 零外部依赖 |
> [!warning] 可接受的任务丢失
>
> AI 生成任务可以重新提交,进程重启丢失排队中的任务是可接受的权衡。但如果你的业务场景中 **任务不可重放**(比如支付指令),则必须选择 RabbitMQ 等持久化实现。
### RabbitMQQueue 实现
基于 AMQP 的持久化消息队列,支持手动 ACK 和重试。详见 [[04-RabbitMQ集成]]。
**RabbitMQQueue 特性**:
| 特性 | 行为 |
|------|------|
| 持久化 | 消息 `DeliveryMode=Persistent`,队列 `durable=true` |
| 重试 | 失败 + retry < max → NACK+requeue;超过 → ACK 丢弃 |
| 背压 | 发布失败时返回错误(不阻塞) |
| 依赖 | 需要 RabbitMQ 服务 |
### 双实现对比
| 维度 | MemoryQueue | RabbitMQQueue |
|------|-------------|---------------|
| **依赖** | 无 | RabbitMQ 服务 |
| **持久化** | 无 | 有(磁盘持久化) |
| **重试** | 无 | NACK + requeue |
| **背压** | 阻塞等待 | 返回错误 |
| **部署** | 单机 | 多机 |
| **适用** | 开发/演示 | 生产环境 |
| **消息丢失** | 进程重启丢失 | 服务重启不丢失 |
### Consumer 桥接
Consumer 从 TaskQueue 消费消息,提交到 WorkerPool 执行,实现队列与并发控制的解耦:
```go
func (c *Consumer) Start(ctx context.Context) error {
return c.queue.Consume(ctx, func(msg TaskMessage) error {
return c.pool.Submit(workerpool.Task{
ID: msg.TaskID,
UserID: msg.UserID,
Fn: func(taskCtx context.Context) error {
c.handler(taskCtx, msg) // 执行 Pipeline
return nil
},
})
})
}
```
```mermaid
graph LR
TQ["TaskQueue全局排队"] --> CONSUMER["Consumer桥接"]
CONSUMER --> WP["WorkerPool单机并发"]
WP --> PIPELINE["Pipeline业务逻辑"]
```
### 三级降级策略
Gen2D 在 `cmd/main.go` 中实现了三级降级链:
| 优先级 | 组件 | 条件 | 行为 |
|--------|------|------|------|
| 1 | TaskQueue + Consumer | 初始化成功 | 队列排队 → Consumer 消费 → WorkerPool 执行 |
| 2 | WorkerPool(直接) | TaskQueue 初始化失败 | 直接提交到协程池 |
| 3 | Legacy FIFO Queue | 均不可用 | 串行队列兜底 |
```go
if taskQueue != nil {
// 优先:通过队列提交
taskQueue.Submit(ctx, msg)
} else if workerPool != nil {
// 回退:直接提交到协程池
workerPool.Submit(task)
} else {
// 兜底:旧串行队列
generateQueue.Enqueue(job)
}
```
### Metrics 指标
| Prometheus 指标 | 类型 | Label | 说明 |
|-----------------|------|-------|------|
| `gen2d_queue_depth` | Gauge | `driver` | 当前队列积压深度 |
| `gen2d_queue_submitted_total` | Counter | `driver` | 入队总量 |
| `gen2d_queue_consumed_total` | Counter | `driver` | 出队总量 |
| `gen2d_queue_submit_duration_seconds` | Histogram | `driver` | 入队耗时 |
| `gen2d_queue_errors_total` | Counter | `driver`, `error_type` | 队列错误总量 |
## 关联文档
- [[00-索引]] — 文档导航与架构总览图
- [[01-系统总览]] — 分层架构与依赖注入
- [[02-协程池]] — 有界并发与 per-user 限流
- [[04-RabbitMQ集成]] — 持久化消息与重试机制