消息队列幂等消费与消息去重实战:基于 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>
<!-- Spring Boot Starter Web -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
<version>3.2.5</version>
</dependency>
<!-- Spring Boot Starter AMQP (RabbitMQ) -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
<version>3.2.5</version>
</dependency>
<!-- Spring Boot Starter Data Redis -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
<version>3.2.5</version>
</dependency>
<!-- MySQL Connector -->
<dependency>
<groupId>com.mysql</groupId>
<artifactId>mysql-connector-j</artifactId>
<version>8.3.0</version>
</dependency>
<!-- Spring Boot Starter JDBC -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jdbc</artifactId>
<version>3.2.5</version>
</dependency>
<!-- Lombok -->
<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:
# 手动 ack
acknowledge-mode: manual
# 消费端限流:每次只拉取 10 条
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`),
-- 核心:order_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";

/**
* 延迟队列:消息先进入此队列,TTL 到期后转发到死信交换机
*/
@Bean
public Queue orderDelayQueue() {
Map<String, Object> args = new HashMap<>();
// 设置 TTL(30 分钟 = 1800000ms)
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();
}

/**
* 延迟交换机(普通 direct 交换机,消息带 TTL 进入延迟队列)
*/
@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;

/** 消息唯一 ID(由 Producer 生成,用于追踪) */
private String messageId;

/** 关单原因,如 TIMEOUT、USER_CANCEL 等 */
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);
// 生产环境:此处应记录到 DB 或告警,并进行补偿
}
}
});
}

/**
* 发布延迟关单消息
* @param orderId 订单号
* @param delayMillis 延迟时间(毫秒)
*/
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 -> {
// 在消息属性上设置 TTL(也可以不设置,使用队列的 x-message-ttl)
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;

/**
* Redis 幂等控制工具
* 使用 SETNX + TTL 实现“同一订单同一操作只执行一次”
*/
@Component
public class RedisIdempotenceUtil {

private static final String KEY_PREFIX = "order:close:idempotence:";

private final StringRedisTemplate redisTemplate;

public RedisIdempotenceUtil(StringRedisTemplate redisTemplate) {
this.redisTemplate = redisTemplate;
}

/**
* 尝试获取幂等锁
* @param orderId 订单号
* @param ttlSeconds 锁的过期时间(秒),建议 > 消息最大延迟时间
* @return true-获取成功(首次处理),false-已处理过
*/
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);
}

/**
* 获取当前锁的 value(用于调试)
*/
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; // 1 小时
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 {
// ============ 第一层防线:Redis SETNX 快速去重 ============
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 {
// 业务失败:释放 Redis 锁,nack 重新入队,等待重试
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());
// 异常时释放锁,nack 重新入队
redisIdempotenceUtil.release(orderId);
channel.basicNack(deliveryTag, false, true);
}
}

/**
* 核心关单业务逻辑
* 包含第二层防线:MySQL 唯一约束兜底
*/
private boolean processOrderClose(OrderCloseMessage message) {
String orderId = message.getOrderId();
String messageId = message.getMessageId();
String reason = message.getReason();

System.out.println("[Business] 开始处理关单,orderId=" + orderId);

// === 第二层防线:MySQL 唯一约束兜底 ===
// 先尝试插入关单流水,利用 order_id 唯一约束保证只有一个请求能成功
// 如果插入成功,说明是第一次处理,继续关单
// 如果插入冲突,说明已有关单流水,直接返回“已处理”

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; // 已处理,返回成功让消费者 ack
}
System.out.println("[Business] 关单流水插入成功,orderId=" + orderId);
} catch (org.springframework.dao.DuplicateKeyException e) {
// 唯一约束冲突:说明已经处理过这个订单了
System.out.println("[Business] 唯一约束冲突,orderId=" + orderId + " 已关单,跳过重复处理");
return true; // 返回成功,让消费者 ack 掉重复消息
}

// === 执行真正的关单逻辑 ===
// 1. 更新订单状态为已取消
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);
}

// 2. 释放库存(这里用日志模拟,实际项目中调用库存服务)
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);

// 发布 30 分钟延迟关单消息(此处为了演示,设置为 30 秒)
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 测试场景一:正常延迟关单

步骤:

  1. 调用 /api/order/create 创建订单
  2. 等待 30 秒(预设的延迟时间)
  3. 观察消费者日志和数据库状态

预期结果:

  • 消费者收到关单消息
  • 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 测试场景二:重复消息幂等验证

步骤:

  1. 正常创建订单并等待关单完成
  2. 调用 /api/order/resend-close-message?orderId=xxx 手动发送重复关单消息
  3. 观察消费者日志,验证幂等效果

预期结果:

  • 重复消息到达消费者
  • 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 唯一约束能否兜底

步骤:

  1. 正常创建订单并等待关单完成
  2. 手动删除 Redis 中的幂等 key(模拟 Redis 失效)
  3. 再次发送重复关单消息
  4. 观察消费者日志和数据库

删除 Redis key 命令:

1
2
3
4
5
# 查看 key
redis-cli KEYS "order:close:idempotence:*"

# 删除 key
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 测试场景四:并发消费验证

场景说明: 模拟两个消费者实例同时消费同一条消息(极端情况)。

测试方法:

  1. 启动两个应用实例(端口 8080 和 8081)
  2. 创建订单,等待延迟消息到达
  3. 由于 RabbitMQ 的 prefetch=10 和手动 ack 机制,正常情况下一台实例消费后另一台不会重复消费
  4. 但我们通过两个实例分别发送相同的关单消息(模拟异常场景)来验证并发幂等

实际测试代码:

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