Files
cs-note/hzh/Gen2D/11-Consumer-Producer桥接.md

4.9 KiB

tags, create time
tags create time
consumer-producer
bridge-pattern
decoupling
go
message-queue
2026-06-03 10:50

11. Consumer-Producer 桥接模式

概述

Consumer 结构体桥接 TaskQueue 和 WorkerPool,实现生产者与消费者的彻底解耦。

正文

架构总览

flowchart LR
    subgraph Producer["Producer"]
        A["Generate Handler"]
    end

    subgraph Queue["TaskQueue 接口"]
        B[("Memory Queue")]
        C[("RabbitMQ Queue")]
    end

    subgraph Bridge["Consumer 桥接层"]
        D{{"Consumer"}}
    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 是整个任务调度体系的桥梁,它只做一件事:从队列取消息,提交到协程池。

// 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 层注入实际业务逻辑

工作流程

启动消费循环

sequenceDiagram
    participant Main as main.go
    participant Consumer as 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 管线

解耦的三层设计

flowchart TB
    subgraph "消息源"
        Q["TaskQueue 接口\n(Memory / RabbitMQ)"]
    end
    subgraph "桥接"
        C["Consumer\n(只关心 消费→提交)"]
    end
    subgraph "执行引擎"
        P["WorkerPool\n(只关心 并发控制)"]
    end
    subgraph "业务逻辑"
        H["TaskHandler 回调\n(RunFromTaskMessage)"]
    end

    Q --> C --> P --> H
组件 知道什么 不知道什么
TaskQueue 消息的存储与投递 WorkerPool 的存在
WorkerPool 任务的并发执行 消息来自哪个队列
Consumer 如何桥接两者 具体的业务逻辑
TaskHandler 生成管线的执行 消息来自内存还是 RabbitMQ

设计哲学:每个组件只关心自己的职责边界,可独立替换、测试、扩展。


Handler 层注入

Consumer 不硬编码业务逻辑,而是通过 TaskHandler 函数签名由外部注入:

// 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 管线。


信号处理与优雅关闭

sequenceDiagram
    participant OS as OS
    participant Main as main.go
    participant Consumer as 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 / 内存资源

关联文档