--- 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
buffered"] MEM_SUBMIT["阻塞写入"] MEM_CONSUME["持续消费"] end subgraph RabbitMQ["RabbitMQQueue"] RMQ_PUB["channel.Publish
Persistent"] RMQ_CONSUME["channel.Consume
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集成]] — 持久化消息与重试机制