消息队列幂等消费与消息去重实战:基于 MySQL 订单系统的 Redis + RabbitMQ 方案
一、问题场景:一条延迟消息引发的“血案”
订单系统里有一个经典场景:用户下单后 30 分钟未支付,系统自动关单并释放库存。
我们采用 RabbitMQ 延迟队列(通过死信交换机 + TTL 实现)来触发关单逻辑。架构如下:
1
| 订单服务 ──发布延迟消息──> RabbitMQ(TTL 30min)──> 死信交换机 ──> 关单队列 ──> 关单消费者
|
这个方案本身没有问题,但上线后我们遇到了三类典型的重复消费场景:
| 场景 |
原因 |
后果 |
| 消息重复投递 |
网络抖动导致 Producer 重试,RabbitMQ Broker 未收到 ack 但消息已入队 |
同一条消息被投递两次 |
| 消费者重试 |
关单逻辑执行成功但 ack 超时,Broker 重新投递 |
关单操作执行两次 |
| 并发消费 |
消费者实例水平扩容,同一消息被不同实例同时拉取(极端情况) |
库存重复释放 |
这些问题的本质是:消息中间件保证的是“至少一次投递”(At-Least-Once),而不保证“恰好一次”(Exactly-Once)。
个人感悟:做分布式系统久了会发现,“恰好一次”在工程上是一个伪命题。你不可能真正避免重复投递,只能让业务逻辑对重复消息“免疫”。想通这一点后,很多看似复杂的中间件问题都变得清晰了——不要试图改变消息中间件的语义,而是让消费端具备幂等能力。
二、方案设计:三层防线
我们设计了三层防线来保证关单操作的幂等性,从快到慢、从轻到重:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31
| ┌─────────────────────────────────────────────────────────┐ │ 消息到达消费者 │ └────────────────────────┬────────────────────────────────┘ ▼ ┌───────────────────────────────┐ │ 第一层:Redis SETNX 快速去重 │ │ 判断 order_id 是否已处理 │ └───────────────┬───────────────┘ │ ┌─────────┴─────────┐ │ 未处理(可继续) │ 已处理(直接 ack 丢弃) ▼ ▼ ┌───────────────┐ ┌────────────┐ │ 执行关单业务 │ │ 直接 ack │ │ 更新订单状态 │ │ 返回 │ │ 释放库存 │ └────────────┘ └───────┬───────┘ │ ▼ ┌───────────────────────────────┐ │ 第二层:MySQL 唯一约束兜底 │ │ order_id 唯一,插入/更新冲突 │ │ 则说明已处理 │ └───────┬───────────────────────┘ │ ▼ ┌───────────────────────────────┐ │ 第三层:手动 ack + 死信队列 │ │ 业务失败时 nack 进入重试 │ │ 超过最大重试次数进入死信 │ └───────────────────────────────┘
|
2.1 为什么用 Redis SETNX 做第一层?
- 快:Redis 单机 QPS 可达 10w+,对消息消费的延迟影响极小
- 简单:
SET key value NX EX ttl 一条命令完成“不存在才设置”的原子操作
- 可过期:设置合理的 TTL,避免 Redis 内存被历史订单 ID 撑爆
2.2 为什么还需要 MySQL 唯一约束兜底?
Redis 可能因为以下原因失效:
- Redis 重启且未开启持久化,key 全部丢失
- Redis 主从切换,主节点刚写入的 key 未同步到从节点
- TTL 设置过短,消息延迟超过 TTL 后重复消息又来了
所以 MySQL 唯一约束是最终的、不可绕过的一道防线。
2.3 手动 ack + 死信队列解决什么?
当一个消息的业务处理持续失败(比如依赖的下游服务宕机),不能无限重试。我们需要:
- 消费端手动 ack,只有在业务执行成功后才确认消息
- 失败时使用
basicNack 将消息重新入队,设置有限次数的重试
- 超过最大重试次数后,将消息转入死信队列,由人工介入处理
三、环境准备
在开始写代码之前,我们先准备好基础环境。
3.1 依赖配置(Maven pom.xml)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39
| <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> <version>3.2.5</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> <version>3.2.5</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-redis</artifactId> <version>3.2.5</version> </dependency> <dependency> <groupId>com.mysql</groupId> <artifactId>mysql-connector-j</artifactId> <version>8.3.0</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-jdbc</artifactId> <version>3.2.5</version> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <version>1.18.32</version> <scope>provided</scope> </dependency> </dependencies>
|
3.2 配置文件(application.yml)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44
| server: port: 8080
spring: datasource: url: jdbc:mysql://localhost:3306/order_db?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true username: root password: root123 driver-class-name: com.mysql.cj.jdbc.Driver hikari: maximum-pool-size: 20 minimum-idle: 5
data: redis: host: localhost port: 6379 database: 0 lettuce: pool: max-active: 50 max-idle: 20 min-idle: 5
rabbitmq: host: localhost port: 5672 username: guest password: guest publisher-confirm-type: correlated publisher-returns: true listener: simple: acknowledge-mode: manual prefetch: 10 retry: enabled: true max-attempts: 3 initial-interval: 2000ms
|
3.3 MySQL 表结构
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25
| CREATE TABLE `t_order` ( `id` BIGINT NOT NULL AUTO_INCREMENT COMMENT '主键', `order_id` VARCHAR(64) NOT NULL COMMENT '订单号', `user_id` BIGINT NOT NULL COMMENT '用户 ID', `status` TINYINT NOT NULL DEFAULT 0 COMMENT '订单状态:0-待支付,1-已支付,2-已取消', `amount` DECIMAL(10, 2) NOT NULL COMMENT '订单金额', `create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间', `update_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间', PRIMARY KEY (`id`), UNIQUE KEY `uk_order_id` (`order_id`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='订单表';
CREATE TABLE `t_order_close_log` ( `id` BIGINT NOT NULL AUTO_INCREMENT COMMENT '主键', `order_id` VARCHAR(64) NOT NULL COMMENT '订单号', `close_reason` VARCHAR(255) NOT NULL DEFAULT 'TIMEOUT' COMMENT '关单原因', `message_id` VARCHAR(128) NOT NULL COMMENT '消息唯一 ID', `create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间', PRIMARY KEY (`id`), UNIQUE KEY `uk_order_id` (`order_id`), KEY `idx_message_id` (`message_id`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='关单流水表';
|
四、代码实现
4.1 RabbitMQ 配置类
首先声明交换机、队列和绑定关系。我们使用死信交换机来实现延迟队列:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105
| package com.example.order.config;
import org.springframework.amqp.core.*; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration;
import java.util.HashMap; import java.util.Map;
@Configuration public class RabbitMQConfig {
public static final String ORDER_DELAY_EXCHANGE = "order.delay.exchange"; public static final String ORDER_DELAY_QUEUE = "order.delay.queue"; public static final String ORDER_DLX_EXCHANGE = "order.dlx.exchange"; public static final String ORDER_CLOSE_QUEUE = "order.close.queue"; public static final String ORDER_CLOSE_DEAD_QUEUE = "order.close.dead.queue"; public static final String ORDER_CLOSE_ROUTING_KEY = "order.close";
@Bean public Queue orderDelayQueue() { Map<String, Object> args = new HashMap<>(); args.put("x-message-ttl", 30 * 60 * 1000); args.put("x-dead-letter-exchange", ORDER_DLX_EXCHANGE); args.put("x-dead-letter-routing-key", ORDER_CLOSE_ROUTING_KEY); return QueueBuilder.durable(ORDER_DELAY_QUEUE) .withArguments(args) .build(); }
@Bean public DirectExchange orderDelayExchange() { return new DirectExchange(ORDER_DELAY_EXCHANGE); }
@Bean public DirectExchange orderDlxExchange() { return new DirectExchange(ORDER_DLX_EXCHANGE); }
@Bean public Queue orderCloseQueue() { Map<String, Object> args = new HashMap<>(); args.put("x-dead-letter-exchange", ORDER_DLX_EXCHANGE); args.put("x-dead-letter-routing-key", "order.close.dead"); return QueueBuilder.durable(ORDER_CLOSE_QUEUE) .withArguments(args) .build(); }
@Bean public Queue orderCloseDeadQueue() { return QueueBuilder.durable(ORDER_CLOSE_DEAD_QUEUE).build(); }
@Bean public Binding delayBinding() { return BindingBuilder.bind(orderDelayQueue()) .to(orderDelayExchange()) .with("order.delay"); }
@Bean public Binding closeQueueBinding() { return BindingBuilder.bind(orderCloseQueue()) .to(orderDlxExchange()) .with(ORDER_CLOSE_ROUTING_KEY); }
@Bean public Binding closeDeadQueueBinding() { return BindingBuilder.bind(orderCloseDeadQueue()) .to(orderDlxExchange()) .with("order.close.dead"); } }
|
4.2 消息体定义
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33
| package com.example.order.message;
import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor;
import java.io.Serializable;
@Data @NoArgsConstructor @AllArgsConstructor public class OrderCloseMessage implements Serializable {
private static final long serialVersionUID = 1L;
private String orderId;
private String messageId;
private String reason;
@Override public String toString() { return "OrderCloseMessage{" + "orderId='" + orderId + '\'' + ", messageId='" + messageId + '\'' + ", reason='" + reason + '\'' + '}'; } }
|
4.3 消息生产者(发布延迟关单消息)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57
| package com.example.order.producer;
import com.example.order.config.RabbitMQConfig; import com.example.order.message.OrderCloseMessage; import org.springframework.amqp.rabbit.connection.CorrelationData; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.stereotype.Component;
import java.util.UUID;
@Component public class OrderCloseProducer {
private final RabbitTemplate rabbitTemplate;
public OrderCloseProducer(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; this.rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> { if (correlationData != null) { String messageId = correlationData.getId(); if (ack) { System.out.println("[Producer] 消息发送成功,messageId=" + messageId); } else { System.out.println("[Producer] 消息发送失败,messageId=" + messageId + ",原因=" + cause); } } }); }
public void publishOrderCloseMessage(String orderId, long delayMillis) { String messageId = UUID.randomUUID().toString(); OrderCloseMessage message = new OrderCloseMessage(orderId, messageId, "TIMEOUT");
CorrelationData correlationData = new CorrelationData(messageId);
rabbitTemplate.convertAndSend( RabbitMQConfig.ORDER_DELAY_EXCHANGE, "order.delay", message, processedMessage -> { processedMessage.getMessageProperties().setExpiration(String.valueOf(delayMillis)); return processedMessage; }, correlationData );
System.out.println("[Producer] 发布延迟关单消息,orderId=" + orderId + ",delay=" + delayMillis + "ms"); } }
|
4.4 Redis 幂等工具类
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60
| package com.example.order.util;
import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.stereotype.Component;
import java.time.Duration;
@Component public class RedisIdempotenceUtil {
private static final String KEY_PREFIX = "order:close:idempotence:";
private final StringRedisTemplate redisTemplate;
public RedisIdempotenceUtil(StringRedisTemplate redisTemplate) { this.redisTemplate = redisTemplate; }
public boolean tryAcquire(String orderId, long ttlSeconds) { String key = KEY_PREFIX + orderId; Boolean success = redisTemplate.opsForValue() .setIfAbsent(key, String.valueOf(System.currentTimeMillis()), Duration.ofSeconds(ttlSeconds)); return Boolean.TRUE.equals(success); }
public boolean isProcessed(String orderId) { String key = KEY_PREFIX + orderId; return Boolean.TRUE.equals(redisTemplate.hasKey(key)); }
public void release(String orderId) { String key = KEY_PREFIX + orderId; redisTemplate.delete(key); System.out.println("[Idempotence] 释放幂等锁,orderId=" + orderId); }
public String getValue(String orderId) { String key = KEY_PREFIX + orderId; return redisTemplate.opsForValue().get(key); } }
|
4.5 关单消费者(核心逻辑)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120
| package com.example.order.consumer;
import com.example.order.message.OrderCloseMessage; import com.example.order.util.RedisIdempotenceUtil; import com.rabbitmq.client.Channel; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.amqp.support.AmqpHeaders; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.messaging.handler.annotation.Header; import org.springframework.stereotype.Component;
import java.io.IOException; import java.time.LocalDateTime; import java.time.format.DateTimeFormatter;
@Component public class OrderCloseConsumer {
private static final int IDEMPOTENCE_TTL_SECONDS = 3600; private static final int MAX_RETRY_COUNT = 3;
private final RedisIdempotenceUtil redisIdempotenceUtil; private final JdbcTemplate jdbcTemplate;
public OrderCloseConsumer(RedisIdempotenceUtil redisIdempotenceUtil, JdbcTemplate jdbcTemplate) { this.redisIdempotenceUtil = redisIdempotenceUtil; this.jdbcTemplate = jdbcTemplate; }
@RabbitListener(queues = "order.close.queue") public void onMessage(OrderCloseMessage message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException {
String orderId = message.getOrderId(); String messageId = message.getMessageId();
System.out.println("[" + Thread.currentThread().getName() + "] 收到关单消息,orderId=" + orderId + ",messageId=" + messageId + ",deliveryTag=" + deliveryTag);
try { boolean acquired = redisIdempotenceUtil.tryAcquire(orderId, IDEMPOTENCE_TTL_SECONDS); if (!acquired) { System.out.println("[Idempotence] orderId=" + orderId + " 已在 Redis 中标记为已处理,直接 ack 丢弃"); channel.basicAck(deliveryTag, false); return; }
boolean businessSuccess = processOrderClose(message);
if (businessSuccess) { System.out.println("[Consumer] 关单业务执行成功,orderId=" + orderId + ",手动 ack"); channel.basicAck(deliveryTag, false); } else { redisIdempotenceUtil.release(orderId); System.out.println("[Consumer] 关单业务执行失败,orderId=" + orderId + ",nack 重新入队"); channel.basicNack(deliveryTag, false, true); }
} catch (Exception e) { System.err.println("[Consumer] 处理消息异常,orderId=" + orderId + ",异常=" + e.getMessage()); redisIdempotenceUtil.release(orderId); channel.basicNack(deliveryTag, false, true); } }
private boolean processOrderClose(OrderCloseMessage message) { String orderId = message.getOrderId(); String messageId = message.getMessageId(); String reason = message.getReason();
System.out.println("[Business] 开始处理关单,orderId=" + orderId);
try { String insertSql = "INSERT INTO t_order_close_log (order_id, close_reason, message_id) VALUES (?, ?, ?)"; int inserted = jdbcTemplate.update(insertSql, orderId, reason, messageId); if (inserted != 1) { System.out.println("[Business] 关单流水插入失败(已存在),orderId=" + orderId); return true; } System.out.println("[Business] 关单流水插入成功,orderId=" + orderId); } catch (org.springframework.dao.DuplicateKeyException e) { System.out.println("[Business] 唯一约束冲突,orderId=" + orderId + " 已关单,跳过重复处理"); return true; }
String updateSql = "UPDATE t_order SET status = 2 WHERE order_id = ? AND status = 0"; int affected = jdbcTemplate.update(updateSql, orderId); if (affected == 0) { System.out.println("[Business] 订单状态已不是待支付,无需关单,orderId=" + orderId); } else { System.out.println("[Business] 订单已关单,orderId=" + orderId); }
System.out.println("[Business] 释放库存完成,orderId=" + orderId);
System.out.println("[Business] 关单处理完成,orderId=" + orderId + ",时间=" + LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")));
return true; } }
|
4.6 订单服务入口(用于测试)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45
| package com.example.order.controller;
import com.example.order.producer.OrderCloseProducer; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.web.bind.annotation.*;
import java.util.UUID;
@RestController @RequestMapping("/api/order") public class OrderController {
private final JdbcTemplate jdbcTemplate; private final OrderCloseProducer orderCloseProducer;
public OrderController(JdbcTemplate jdbcTemplate, OrderCloseProducer orderCloseProducer) { this.jdbcTemplate = jdbcTemplate; this.orderCloseProducer = orderCloseProducer; }
@PostMapping("/create") public String createOrder(@RequestParam String userId, @RequestParam String amount) { String orderId = UUID.randomUUID().toString();
String sql = "INSERT INTO t_order (order_id, user_id, status, amount) VALUES (?, ?, 0, ?)"; jdbcTemplate.update(sql, orderId, userId, amount);
orderCloseProducer.publishOrderCloseMessage(orderId, 30 * 1000);
return "订单创建成功,orderId=" + orderId; }
@PostMapping("/resend-close-message") public String resendCloseMessage(@RequestParam String orderId) { orderCloseProducer.publishOrderCloseMessage(orderId, 0); return "重复关单消息已发送,orderId=" + orderId; } }
|
五、并发测试与验证
5.1 测试环境
| 组件 |
版本 |
说明 |
| JDK |
17.0.10 |
编译和运行环境 |
| Spring Boot |
3.2.5 |
应用框架 |
| MySQL |
8.0.36 |
存储订单和关单流水 |
| Redis |
7.2.4 |
幂等标记 |
| RabbitMQ |
3.13.2 |
消息队列 |
5.2 测试场景一:正常延迟关单
步骤:
- 调用
/api/order/create 创建订单
- 等待 30 秒(预设的延迟时间)
- 观察消费者日志和数据库状态
预期结果:
- 消费者收到关单消息
- Redis 中生成
order:close:idempotence:{orderId} key
t_order_close_log 表插入一条关单流水
t_order 表订单状态更新为 2(已取消)
实际日志输出:
1 2 3 4 5 6 7 8 9 10 11 12
| [Producer] 消息发送成功,messageId=8f3a2b1c-... [Producer] 发布延迟关单消息,orderId=5a7d4e9f-...,delay=30000ms
--- 30 秒后 ---
[org.springframework.amqp.rabbit.RabbitListenerEndpointContainer#0-1] 收到关单消息,orderId=5a7d4e9f-...,messageId=8f3a2b1c-...,deliveryTag=1 [Business] 开始处理关单,orderId=5a7d4e9f-... [Business] 关单流水插入成功,orderId=5a7d4e9f-... [Business] 订单已关单,orderId=5a7d4e9f-... [Business] 释放库存完成,orderId=5a7d4e9f-... [Business] 关单处理完成,orderId=5a7d4e9f-...,时间=2026-08-31 08:34:12 [Consumer] 关单业务执行成功,orderId=5a7d4e9f-...,手动 ack
|
5.3 测试场景二:重复消息幂等验证
步骤:
- 正常创建订单并等待关单完成
- 调用
/api/order/resend-close-message?orderId=xxx 手动发送重复关单消息
- 观察消费者日志,验证幂等效果
预期结果:
- 重复消息到达消费者
- Redis
tryAcquire 返回 false
- 消息被直接 ack,不执行业务逻辑
t_order_close_log 表中只有一条流水记录
实际日志输出:
1 2 3 4
| [Producer] 发布延迟关单消息,orderId=5a7d4e9f-...,delay=0ms
[org.springframework.amqp.rabbit.RabbitListenerEndpointContainer#0-1] 收到关单消息,orderId=5a7d4e9f-...,messageId=7c2d9e1a-...,deliveryTag=2 [Idempotence] orderId=5a7d4e9f-... 已在 Redis 中标记为已处理,直接 ack 丢弃
|
数据库验证(MySQL 8.0):
1 2 3 4 5 6 7
| mysql> SELECT * FROM t_order_close_log WHERE order_id = '5a7d4e9f-...'; + | id | order_id | close_reason | message_id | create_time | + | 1 | 5a7d4e9f-... | TIMEOUT | 8f3a2b1c-... | 2026-08-31 08:34:12 | + 1 row in set (0.00 sec)
|
只有一条记录,幂等验证通过。
5.4 测试场景三:Redis 宕机时 MySQL 兜底验证
这是最关键的一个场景:模拟 Redis 不可用(或 key 已过期),验证 MySQL 唯一约束能否兜底。
步骤:
- 正常创建订单并等待关单完成
- 手动删除 Redis 中的幂等 key(模拟 Redis 失效)
- 再次发送重复关单消息
- 观察消费者日志和数据库
删除 Redis key 命令:
1 2 3 4 5
| redis-cli KEYS "order:close:idempotence:*"
redis-cli DEL "order:close:idempotence:5a7d4e9f-..."
|
预期结果:
- Redis
tryAcquire 返回 true(因为 key 被删了)
- 消费者进入业务逻辑
- 尝试插入
t_order_close_log 时触发唯一约束冲突(DuplicateKeyException)
- 捕获异常并返回
true,消息被 ack
- 数据库不会出现重复流水
实际日志输出:
1 2 3 4
| [org.springframework.amqp.rabbit.RabbitListenerEndpointContainer#0-1] 收到关单消息,orderId=5a7d4e9f-...,messageId=9e1b3c5d-...,deliveryTag=3 [Business] 开始处理关单,orderId=5a7d4e9f-... [Business] 唯一约束冲突,orderId=5a7d4e9f-... 已关单,跳过重复处理 [Consumer] 关单业务执行成功,orderId=5a7d4e9f-...,手动 ack
|
数据库验证:
1 2 3 4 5 6 7
| mysql> SELECT COUNT(*) FROM t_order_close_log WHERE order_id = '5a7d4e9f-...'; + | COUNT(*) | + | 1 | + 1 row in set (0.00 sec)
|
仍然只有一条流水,MySQL 唯一约束兜底成功。
5.5 测试场景四:并发消费验证
场景说明: 模拟两个消费者实例同时消费同一条消息(极端情况)。
测试方法:
- 启动两个应用实例(端口 8080 和 8081)
- 创建订单,等待延迟消息到达
- 由于 RabbitMQ 的
prefetch=10 和手动 ack 机制,正常情况下一台实例消费后另一台不会重复消费
- 但我们通过两个实例分别发送相同的关单消息(模拟异常场景)来验证并发幂等
实际测试代码:
package com.example.order;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.jdbc.core.JdbcTemplate;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicInteger;
import static org.junit.jupiter.api.Assertions.assertEquals;
@SpringBootTest
public class ConcurrentIdempotenceTest {
@Autowired
private JdbcTemplate jdbcTemplate;
@Autowired
private com.example.order.producer.OrderCloseProducer producer;
@Test
public void testConcurrentCloseMessages() throws InterruptedException {
String orderId = UUID.randomUUID().toString();
String userId = "1001";
// 创建订单
jdbcTemplate.update(
"INSERT INTO t_order (order_id, user_id, status, amount) VALUES (?, ?, 0, 199.00)",
orderId, userId
);
// 模拟 10 条重复消息同时到达
int threadCount = 10;
ExecutorService executor = Executors.newFixedThreadPool(threadCount);
CountDownLatch latch = new CountDownLatch(threadCount);
AtomicInteger successCount = new AtomicInteger(0);
for (int i = 0; i < threadCount; i++) {
executor.submit(() -> {
try {
producer.publishOrderCloseMessage(orderId, 0);
successCount.incrementAndGet();
} finally {
latch.countDown();
}
});
}
latch.await();
executor.shutdown();
// 等待消费者处理
Thread.sleep(5000);
// 验证:关单流水只有一条
Integer closeLogCount = jdbcTemplate.queryForObject(
"SELECT COUNT(*) FROM t_order_close_log WHERE order_id