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

7.2 KiB
Raw Blame History

03 - 任务队列 (Task Queue)

一句话概括:可插拔任务队列接口,支持 Memory 和 RabbitMQ 双实现,通过工厂模式一行切换。

架构设计

graph TB
    subgraph Producer["生产者"]
        HANDLER["Handler.Generate()"]
    end

    subgraph Interface["TaskQueue 接口"]
        SUBMIT["Submit(ctx, msg)"]
        CONSUME["Consume(ctx, 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 接口

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 聚合错误

工厂模式

通过配置驱动,一行切换队列实现:

tq, err := taskqueue.New(cfg.TaskQueue)
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 消息结构

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 的内存队列,零外部依赖。

核心逻辑:

type MemoryQueue struct {
    ch      chan TaskMessage  // 有界缓冲
    closed  atomic.Bool
    closeCh chan struct{}
}

Submit

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
    }
}

💡 阻塞语义:MemoryQueue 的 Submit 是阻塞的——当 buffer 满时,调用方会阻塞直到有空位或 ctx 取消。这与 WorkerPool 的非阻塞拒绝形成对比。

Consume

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 错误直接丢弃)
背压 阻塞直到有空位
依赖 零外部依赖

⚠️ 可接受的任务丢失:AI 生成任务可以重新提交,进程重启丢失排队中的任务是可接受的权衡。

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 执行,实现队列与并发控制的解耦:

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
            },
        })
    })
}
TaskQueue → Consumer → WorkerPool → Pipeline
   全局排队      桥接       单机并发     业务逻辑

三级降级策略

Gen2D 在 cmd/main.go 中实现了三级降级链:

优先级 组件 条件 行为
1 TaskQueue + Consumer 初始化成功 队列排队 → Consumer 消费 → WorkerPool 执行
2 WorkerPool(直接) TaskQueue 初始化失败 直接提交到协程池
3 Legacy FIFO Queue 均不可用 串行队列兜底
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 队列错误总量

关联文档