☰
【RabbitMQ #12】 | 业务幂等
2026/10/1 15:41:38 网站建设 项目流程

一、什么是幂等

幂等:同一个业务,执行一次或者执行多次,对业务最终状态产生的影响完全一致。 典型场景:表单重复提交、MQ 消息重复投递,多次执行不会产生脏数据。

网页表单场景:进入页面下发 token,提交时校验 token,提交成功后删除 token,利用令牌机制实现表单幂等。

MQ 场景为什么会重复消费? MQ 消费者在业务处理完成前宕机,ACK 没有发送成功,MQ 会重新投递消息,造成重复消费,因此消费端必须做幂等。

二、RabbitMQ 实现幂等的两种方案

方案 1:唯一消息 ID(消息 ID 幂等)

核心思路:给每条消息分配全局唯一 ID,消费成功后存入幂等表;再次收到相同 ID 则判定为重复消息,直接跳过业务。

流程: ① 生产者生成唯一消息 id,随消息一起发送给消费者 ② 消费者收到消息,先查询幂等表; ③ 如果 ID 已存在 → 重复消息,直接 ACK,不执行业务; ④ 如果 ID 不存在 → 执行业务,业务成功后写入幂等记录(业务入库与幂等记录入库建议放在同一个事务,保证原子性)

幂等表 SQL

CREATE TABLE msg_id_record ( id BIGINT PRIMARY KEY AUTO_INCREMENT, msg_id VARCHAR(64) NOT NULL COMMENT '消息唯一ID', create_time DATETIME DEFAULT CURRENT_TIMESTAMP, UNIQUE uk_msg_id (msg_id) -- 唯一索引兜底,防止重复插入 ) ENGINE=InnoDB;
Go 生产者代码(发送消息携带 messageId)
package main import ( "github.com/rabbitmq/amqp091-go" "github.com/google/uuid" // go get github.com/google/uuid "log" ) func main() { conn, err := amqp091.Dial("amqp://guest:guest@127.0.0.1:5672/") if err != nil { log.Fatal(err) } defer conn.Close() ch, err := conn.Channel() if err != nil { log.Fatal(err) } defer ch.Close() // 生成全局唯一消息ID,等价Spring自动创建的messageId msgId := uuid.NewString() err = ch.Publish( "", "test_queue", false, false, amqp091.Publishing{ ContentType: "application/json", MessageId: msgId, // MQ原生MessageId字段,推荐 Body: []byte(`{"orderId":1001}`), }) if err != nil { log.Fatal(err) } log.Printf("消息发送成功, messageId=%s", msgId) }
Go 消费者代码(消费时做幂等判断)

核心逻辑:收到消息,先查幂等表,如果 messageId 已存在,直接跳过业务,返回 ack;不存在,执行业务 + 写入幂等记录(业务和幂等记录同事务,保证原子性)

package main import ( "github.com/rabbitmq/amqp091-go" "log" ) // 查询数据库,判断该消息ID是否已经处理过 func isMsgHandled(msgId string) (bool, error) { // SELECT count(*) FROM msg_id_record WHERE msg_id = ? return false, nil } // 插入消息ID到幂等记录表 func saveMsgId(msgId string) error { // INSERT INTO msg_id_record(msg_id,create_time) VALUES (?,now()) return nil } func handleBusiness(body []byte) error { // 业务逻辑 return nil } func main() { conn, err := amqp091.Dial("amqp://guest:guest@127.0.0.1:5672/") if err != nil { log.Fatal(err) } defer conn.Close() ch, err := conn.Channel() if err != nil { log.Fatal(err) } defer ch.Close() msgs, err := ch.Consume( "test_queue", "", false, // autoAck=false 手动ack false, false, false, nil, ) if err != nil { log.Fatal(err) } forever := make(chan struct{}) go func() { for d := range msgs { msgId := d.MessageId log.Printf("收到消息 messageId=%s", msgId) // 1.幂等校验 handled, err := isMsgHandled(msgId) if err != nil { log.Println("查询幂等表失败", err) _ = d.Nack(false, true) // 查询异常,消息重投 continue } if handled { log.Println("重复消息,直接ack丢弃,不执行业务") _ = d.Ack(false) continue } // 2.执行业务逻辑(业务和saveMsgId尽量在同一个数据库事务) err = handleBusiness(d.Body) if err != nil { log.Println("业务处理失败", err) _ = d.Nack(false, true) continue } // 3.业务成功,保存消息ID到幂等表 err = saveMsgId(msgId) if err != nil { log.Println("保存消息ID失败", err) _ = d.Nack(false, true) continue } // 4.全部成功,ack _ = d.Ack(false) } }() log.Println("消费者启动") <-forever }

✅ 优点:通用性强,任何 MQ 消息场景都能用 ❌ 缺点:需要单独维护一张幂等记录表,多一次数据库查询,有额外性能开销


方案 2:基于业务状态做判断(推荐,无需额外幂等表)

核心思路:利用业务自身状态流转做判断,更新 SQL 带上状态条件,数据库原子执行,不需要额外幂等表。

场景示例:支付成功后修改订单状态,订单状态:1 = 未支付,2 = 已支付 逻辑:只有订单状态为未支付时,才允许更新为已支付;如果订单已经是已支付,直接忽略。

不要先查询再更新(存在并发间隙,会有并发安全问题),直接把状态判断写进 UPDATE 的 WHERE 条件,数据库原子完成判断 + 更新。

UPDATE order SET status = 2, pay_time=NOW() WHERE id = ? AND status = 1

执行后判断RowsAffected,返回 0 代表:订单不存在 / 状态不是未支付(重复消息,无需处理)

Go + GORM 实现
package main import ( "github.com/rabbitmq/amqp091-go" "gorm.io/gorm" "log" "time" ) // Order 订单实体 type Order struct { ID uint64 `gorm:"primaryKey"` Status int // 1:未支付 2:已支付 PayTime time.Time `gorm:"column:pay_time"` } // 消费函数,支付成功后更新订单状态 func listenOrderPay(orderId uint64, db *gorm.DB) error { // 原子更新:仅当status=1时才更新为已支付,带上payTime result := db.Model(&Order{}). Where("id = ? AND status = ?", orderId, 1). Updates(map[string]any{ "status": 2, "pay_time": time.Now(), }) if result.Error != nil { return result.Error } // RowsAffected == 0 代表:订单不存在 / 状态不是未支付(已经处理过,重复消息) if result.RowsAffected == 0 { log.Printf("无需更新,订单已处理或不存在 orderId=%d", orderId) return nil } log.Printf("订单更新为已支付成功 orderId=%d", orderId) return nil } func main() { // 1. 连接rabbitmq(省略db初始化代码) conn, err := amqp091.Dial("amqp://guest:guest@127.0.0.1:5672/") if err != nil { log.Fatal(err) } defer conn.Close() ch, err := conn.Channel() if err != nil { log.Fatal(err) } defer ch.Close() // 监听队列 mark.order.pay.queue,topic交换机 pay.topic,routingKey pay.success msgs, err := ch.Consume( "mark.order.pay.queue", "", false, // 手动ack false, false, false, nil, ) if err != nil { log.Fatal(err) } // 消费协程 go func() { for d := range msgs { // 消息体是orderId,实际项目需要json.Unmarshal解析消息体 var orderId uint64 // json.Unmarshal(d.Body, &orderId) err := listenOrderPay(orderId, db) if err != nil { log.Println("消费失败", err) _ = d.Nack(false, true) // 失败,重新入队 } else { _ = d.Ack(false) // 成功,ack } } }() log.Println("订单支付监听启动") <-make(chan struct{}) }

优点:不用额外幂等表,利用业务本身状态,性能更好,适合订单、库存这类有状态流转场景 ❌ 缺点:依赖业务状态,不是所有业务都适用

三、业务场景:支付服务与交易服务,保证订单状态最终一致性

面试简答背诵版Q:如何保证支付服务与交易服务之间的订单状态一致性?

  • 首先,支付服务会在用户支付成功以后利用 MQ 消息通知交易服务,完成订单状态同步。
  • 其次,为了保证 MQ 消息的可靠性,我们采用了生产者确认机制、消费者确认、消费者失败重试等策略,确保消息投递和处理的可靠性。同时也开启了 MQ 的持久化,避免因服务器宕机导致消息丢失。
  • 最后,我们还在交易服务更新订单状态时做了业务幂等判断,避免因消息重复消费导致订单状态异常。

Q:如果交易服务消息处理失败,兜底方案?

  • 在交易服务设置定时任务,定期查询订单支付状态。即便 MQ 通知失败,定时任务作为兜底,保证订单支付状态的最终一致性。

四、面试小结

  1. MQ 重复消费根源:消费者业务处理完成前宕机,ACK 丢失,MQ 重投消息。
  2. 幂等两种方案:消息唯一 ID(通用,需要幂等表)、业务状态判断(推荐,无额外表,利用数据库原子更新)
  3. 分布式最终一致性:MQ 可靠投递 + 消费端幂等 + 定时任务兜底。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询