【RabbitMQ #12】 | 业务幂等
一、什么是幂等幂等同一个业务执行一次或者执行多次对业务最终状态产生的影响完全一致。 典型场景表单重复提交、MQ 消息重复投递多次执行不会产生脏数据。网页表单场景进入页面下发 token提交时校验 token提交成功后删除 token利用令牌机制实现表单幂等。MQ 场景为什么会重复消费 MQ 消费者在业务处理完成前宕机ACK 没有发送成功MQ 会重新投递消息造成重复消费因此消费端必须做幂等。二、RabbitMQ 实现幂等的两种方案方案 1唯一消息 ID消息 ID 幂等核心思路给每条消息分配全局唯一 ID消费成功后存入幂等表再次收到相同 ID 则判定为重复消息直接跳过业务。流程 ① 生产者生成唯一消息 id随消息一起发送给消费者 ② 消费者收到消息先查询幂等表 ③ 如果 ID 已存在 → 重复消息直接 ACK不执行业务 ④ 如果 ID 不存在 → 执行业务业务成功后写入幂等记录业务入库与幂等记录入库建议放在同一个事务保证原子性幂等表 SQLCREATE 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) -- 唯一索引兜底防止重复插入 ) ENGINEInnoDB;Go 生产者代码发送消息携带 messageIdpackage 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:guest127.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:guest127.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, // autoAckfalse 手动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_timeNOW() 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 { // 原子更新仅当status1时才更新为已支付带上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:guest127.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.queuetopic交换机 pay.topicroutingKey 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 通知失败定时任务作为兜底保证订单支付状态的最终一致性。四、面试小结MQ 重复消费根源消费者业务处理完成前宕机ACK 丢失MQ 重投消息。幂等两种方案消息唯一 ID通用需要幂等表、业务状态判断推荐无额外表利用数据库原子更新分布式最终一致性MQ 可靠投递 消费端幂等 定时任务兜底。
上一篇/下一篇内容由系统自动关联
返回资讯列表 →