--- tags: [rabbitmq, amqp, message-queue, persistence, ack-nack, retry-pattern, go] create time: 2026-06-03 10:15 --- # 04. RabbitMQ 集成 ## 概述 基于 AMQP 协议的持久化消息队列实现,支持手动 ACK/NACK、失败重试和死信丢弃,保障消息不丢失。 > **一句话概括**:基于 AMQP 的持久化消息队列,支持手动 ACK、失败重试和死信丢弃。 ## 正文 ### 消息流 ```mermaid graph TB subgraph Producer["生产者"] HANDLER["Handler.Generate"] end subgraph RabbitMQ["RabbitMQ"] EXCHANGE["Default ExchangeDirect"] QUEUE["gen2d:tasks
durable=true"] end subgraph Consumer["消费者"] CONSUME["channel.Consume
autoAck=false"] end subgraph Decision["ACK/NACK决策树"] SUCCESS{"handler成功?"} RETRY{"retry lt maxRetry?"} ACK_OK["ACK确认消费"] NACK["NACK + requeue
重新入队"] ACK_DISCARD["ACK discard
丢弃死信"] end HANDLER -->|"Publish Persistent"| EXCHANGE EXCHANGE --> QUEUE QUEUE --> CONSUME CONSUME --> SUCCESS SUCCESS -->|是| ACK_OK SUCCESS -->|否| RETRY RETRY -->|是| NACK RETRY -->|否| ACK_DISCARD NACK -.->|"重新投递"| QUEUE ``` ### 连接流程 RabbitMQQueue 在初始化时完成连接、声明队列、设置 QoS: ```mermaid graph TB DIAL["amqp.Dialcfg.URL"] --> CHANNEL["conn.Channel"] CHANNEL --> DECLARE["ch.QueueDeclare name durable=true"] DECLARE --> QOS["ch.Qosprefetch=1"] QOS --> RESULT["RabbitMQQueue实例"] ``` **参数说明**: | 参数 | 默认值 | 说明 | |------|--------|------| | `URL` | `amqp://guest:guest@localhost:5672/` | AMQP 连接地址 | | `Queue` | `gen2d:tasks` | 队列名称 | | `Prefetch` | `1` | 每次预取消息数,1 保证公平调度 | | `MaxRetry` | `3` | 失败最大重试次数 | > [!tip] Prefetch=1 的含义 > > 每个 Consumer 同时只处理 1 条消息,处理完ACK后才接收下一条。这避免了消息堆积在 Consumer 端,配合协程池的并发控制实现精确的任务调度。 ### 消息发布 (Submit) ```go func (q *RabbitMQQueue) Submit(ctx context.Context, msg TaskMessage) error { body, _ := json.Marshal(msg) return q.channel.PublishWithContext(ctx, "", // exchange(默认直连) q.queue, // routing key = queue name false, // mandatory false, // immediate amqp.Publishing{ ContentType: "application/json", DeliveryMode: amqp.Persistent, // 持久化消息 Body: body, Timestamp: time.Now(), Headers: amqp.Table{ "x-retry-count": msg.RetryCount, }, }, ) } ``` **关键配置**: | 属性 | 值 | 说明 | |------|-----|------| | `DeliveryMode` | `Persistent (2)` | 消息写磁盘,RabbitMQ 重启不丢失 | | `ContentType` | `application/json` | JSON 序列化 | | `x-retry-count` | `int` (header) | 当前重试次数,供消费端判断 | ### 消息消费 (Consume) ```go func (q *RabbitMQQueue) Consume(ctx context.Context, handler func(TaskMessage) error) error { deliveries, _ := q.channel.Consume( q.queue, // queue "", // consumer name(自动生成) false, // autoAck = false(手动 ACK) false, // exclusive false, // noLocal false, // noWait nil, // args ) // 持续消费循环... } ``` **手动 ACK 模式**: 手动 ACK 给予消费者完全的控制权——只有当消息被成功处理后才确认,否则可以选择重试或丢弃。 ### ACK/NACK 决策树 ```mermaid graph TB MSG["收到消息"] --> PARSE{"JSON解析成功?"} PARSE -->|否| ACK_DISCARD1["ACK discard
格式错误无法恢复"] PARSE -->|是| HANDLER{"handlermsg执行成功?"} HANDLER -->|是| ACK_OK["ACK确认消费"] HANDLER -->|否| RETRY_CHECK{"msg.RetryCount lt maxRetry?"} RETRY_CHECK -->|是| NACK["NACKrequeuetrue
重新入队等待重试"] RETRY_CHECK -->|否| ACK_DISCARD2["ACK discard
超过最大重试记录死信日志"] ``` | 场景 | 操作 | 说明 | |------|------|------| | handler 成功 | `d.Ackfalse` | 确认消费,消息从队列移除 | | handler 失败 + retry < max | `d.Nackfalse, true` | 拒绝并重新入队,retry count 递增 | | handler 失败 + retry >= max | `d.Ackfalse` + 日志 | 超过最大重试,丢弃可扩展为死信队列 | | JSON 解析失败 | `d.Ackfalse` | 格式错误无法恢复,直接丢弃 | **重试计数传递**: ```go // 发布时:写入 header Headers: amqp.Table{ "x-retry-count": msg.RetryCount, } // 消费时:从 header 读取 if retry, ok := d.Headers["x-retry-count"].(int32); ok { msg.RetryCount = int(retry) } ``` > [!warning] NACK requeue 的行为 > > `Nack(false, true)` 会将消息重新放回队列头部。如果消费者立即再次消费,可能导致"毒消息"反复重试。Gen2D 通过 `maxRetry=3` 限制重试次数,并在超过后 ACK 丢弃来规避此问题。 ### Metrics 指标 | Prometheus 指标 | 类型 | Label | 说明 | |-----------------|------|-------|------| | `gen2d_rabbitmq_connection_status` | Gauge | — | 连接状态(1=connected, 0=disconnected) | | `gen2d_queue_submitted_total` | Counter | `driver=rabbitmq` | 发布消息总量 | | `gen2d_queue_consumed_total` | Counter | `driver=rabbitmq` | 成功消费总量 | | `gen2d_queue_errors_total` | Counter | `driver=rabbitmq`, `error_type` | 错误总量 | | `gen2d_queue_submit_duration_seconds` | Histogram | `driver=rabbitmq` | 发布耗时 | ### 关闭流程 ```go func (q *RabbitMQQueue) Close() error { metrics.RabbitMQConnectionStatus.Set(0) // 标记断开 var errs []error errs = append(errs, q.channel.Close()) // 先关 channel errs = append(errs, q.conn.Close()) // 再关连接 return errors.Join(errs...) } ``` **关闭顺序**:Channel 先于 Connection 关闭,确保所有未确认的消息被释放回队列。 ### 配置参考 ```yaml # config.yaml taskqueue: driver: "rabbitmq" rabbitmq: url: "amqp://guest:guest@localhost:5672/" queue: "gen2d:tasks" prefetch: 1 max_retry: 3 ``` | 配置项 | 环境变量 | 默认值 | 说明 | |--------|----------|--------|------| | `url` | `GEN2D_TASKQUEUE_RABBITMQ_URL` | `amqp://guest:guest@localhost:5672/` | AMQP 连接地址 | | `queue` | `GEN2D_TASKQUEUE_RABBITMQ_QUEUE` | `gen2d:tasks` | 队列名称 | | `prefetch` | `GEN2D_TASKQUEUE_RABBITMQ_PREFETCH` | `1` | 预取消息数 | | `max_retry` | `GEN2D_TASKQUEUE_RABBITMQ_MAX_RETRY` | `3` | 最大重试次数 | ## 关联文档 - [[00-索引]] — 文档导航与架构总览图 - [[03-任务队列]] — 可插拔接口与 MemoryQueue 实现 - [[02-协程池]] — 单机并发控制 - [[01-系统总览]] — 分层架构与配置级联