RabbitMQ 详解和代码实现
·
文章目录
第一部分:基础概念
1.1 什么是 RabbitMQ
RabbitMQ 是一个开源的消息代理(Message Broker)和队列服务器,实现了高级消息队列协议(AMQP 0.9.1)。
- 核心价值:
- 解耦: 系统 A 发送消息后无需等待系统 B 处理,降低耦合度。
- 异步: 耗时操作(如发邮件、生成报表)放入队列,主流程立即返回,提升响应速度。
- 削峰填谷: 突发流量先存入队列,消费者按能力慢慢处理,防止系统崩溃。
1.2 核心术语对照表
| 术语 | 英文 | 解释 | 生活类比 |
|---|---|---|---|
| Broker | Broker | RabbitMQ 服务端实体 | 邮局大楼 |
| Virtual Host | VHost | 虚拟主机,用于逻辑隔离 | 邮局里的不同分区(如国内件区、国际件区) |
| Exchange | Exchange | 交换机,接收消息并路由 | 分拣中心 |
| Queue | Queue | 队列,存储消息的缓冲区 | 具体的邮箱/包裹架 |
| Binding | Binding | 连接交换机和队列的规则 | 分拣规则(如:北京的去 A 区) |
| Routing Key | Routing Key | 路由键,消息携带的标签 | 包裹上的地址标签 |
| Producer | Producer | 消息生产者 | 寄件人 |
| Consumer | Consumer | 消息消费者 | 收件人 |
1.3 核心交互流程图
第二部分:项目搭建与基础配置
2.1 标准项目结构
基于您的实战项目,采用清晰的分层结构:
spring-rabbitmq/
├── src/main/java/org/example/rabbitmq/
│ ├── config/ # 【核心】所有 MQ 组件定义 (Exchange, Queue, Binding)
│ │ ├── DeadLetterRabbitConfig.java # 死信队列配置 (异常兜底)
│ │ ├── DelayRabbitConfig.java # 延迟队列配置 (TTL+DLX)
│ │ ├── DirectRabbitConfig.java # 直连交换机 (点对点)
│ │ ├── FanoutRabbitConfig.java # 扇形交换机 (广播)
│ │ ├── TopicRabbitConfig.java # 主题交换机 (模糊匹配)
│ │ ├── MessageListenerConfig.java # 手动确认监听器容器配置
│ │ └── RabbitConfig.java # 全局 RabbitTemplate 及回调配置
│ ├── constant/ # 常量定义,避免硬编码
│ │ └── RabbitMqConstant.java
│ ├── rabbit/ # 注解式接收者 (@RabbitListener)
│ │ ├── DelayReceiver.java # 延迟消息处理
│ │ ├── DirectReceiver.java # 直连消息处理
│ │ ├── FanoutReceiverA/B/C.java # 广播消息处理
│ │ └── TopicMan/TotalReceiver.java # 主题消息处理
│ ├── receiver/ # 接口式接收者 (ChannelAwareMessageListener)
│ │ └── MyAckReceiver.java # 统一手动 ACK 逻辑
│ └── SpringRabbitmqStartApplication.java
├── src/main/resources/
│ ├── application.yaml # 连接与基础行为配置
│ └── logback-spring.xml # 日志配置
└── pom.xml
2.2 依赖配置 (pom.xml)
确保引入以下核心依赖,特别注意 JSON 序列化支持。
<dependencies>
<!-- Spring Boot RabbitMQ Starter -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<!-- 测试支持 -->
<dependency>
<groupId>org.springframework.amqp</groupId>
<artifactId>spring-rabbit-test</artifactId>
<scope>test</scope>
</dependency>
<!-- JSON 序列化支持 (生产环境必须,避免 JDK 默认序列化的坑) -->
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
</dependency>
<!-- 日志增强 (可选) -->
<dependency>
<groupId>net.logstash.logback</groupId>
<artifactId>logstash-logback-encoder</artifactId>
</dependency>
</dependencies>
2.3 核心配置文件 (application.yaml) 深度解析
spring:
rabbitmq:
host: 192.168.188.153
port: 5672 # AMQP 协议端口 (管理界面是 15672)
username: kiv
password: kiv
virtual-host: /rabbitmqkiv # 虚拟主机,实现多租户隔离
# === 生产者可靠性保障 (防止消息丢失第一道防线) ===
# correlated: 开启发布确认,消息到达交换机后回调 ConfirmCallback
publisher-confirm-type: correlated
# true: 消息到达交换机但无法路由到队列时,触发 ReturnCallback
publisher-returns: true
listener:
simple:
# === 消费者可靠性保障 (防止消息丢失第二道防线) ===
# manual: 必须手动调用 channel.basicAck(),否则消息不会移除
acknowledge-mode: manual
# 重试机制 (可选,建议结合死信队列使用)
retry:
enabled: true
max-attempts: 3
initial-interval: 1000ms
2.4 全局配置类 (RabbitConfig.java)
关键点:配置 JSON 序列化转换器,并注册确认回调。这是生产环境的标配。
@Configuration
public class RabbitConfig {
@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
// 1. 设置 JSON 序列化转换器 (重要!避免乱码和类型转换错误)
rabbitTemplate.setMessageConverter(new Jackson2JsonMessageConverter());
// 2. 开启 ConfirmCallback (消息到达交换机)
rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
if (!ack) {
System.err.println("❌ 消息发送失败,未到达交换机: " +
(correlationData != null ? correlationData.getId() : "null") + ", 原因: " + cause);
// TODO: 记录数据库,启动定时任务重发
} else {
System.out.println("✅ 消息成功到达交换机: " +
(correlationData != null ? correlationData.getId() : "null"));
}
});
// 3. 开启 ReturnsCallback (消息未路由到队列)
rabbitTemplate.setReturnsCallback(returnedMessage -> {
System.err.println("❌ 消息路由失败: " + new String(returnedMessage.getMessage().getBody())
+ ", 交换机: " + returnedMessage.getExchange()
+ ", 路由键: " + returnedMessage.getRoutingKey());
// TODO: 记录死信或报警
});
return rabbitTemplate;
}
}
第三部分:五种交换机模式实战
3.1 直连交换机 (Direct Exchange)
- 机制: 精确匹配。
Routing Key必须与Binding Key完全一致。 - 场景: 点对点通知、特定任务分发。
- 注意: 如果 Routing Key 不匹配且未开启 ReturnCallback,消息会直接丢失。
测试类 (DirectSender.java)
@SpringBootTest
class DirectTests {
@Autowired
RabbitTemplate rabbitTemplate; // 使用RabbitTemplate,这提供了接收/发送等等方法
// 测试发送消息到Direct
@Test
public void sendDirectMessage() throws InterruptedException {
// 循环发送消息
while (true) {
String messageId = String.valueOf(UUID.randomUUID());
String messageData = "test message, hello!";
String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
Map<String, Object> map = new HashMap<>();
map.put("messageId", messageId);
map.put("messageData", messageData);
map.put("createTime", createTime);
// 将消息携带绑定键值:TestDirectRouting 发送到交换机 TestDirectExchange
rabbitTemplate.convertAndSend("TestDirectExchange", "TestDirectRouting", map);
// sleep(50);
}
}
}
配置类 (DirectRabbitConfig.java)
@Configuration
public class DirectRabbitConfig {
@Bean
public Queue TestDirectQueue() {
// durable=true: 队列持久化,RabbitMQ 重启后依然存在
return new Queue(RabbitMqConstant.TEST_DIRECT_QUEUE, true);
}
@Bean
DirectExchange TestDirectExchange() {
return new DirectExchange(RabbitMqConstant.TEST_DIRECT_EXCHANGE, true, false);
}
@Bean
Binding bindingDirect() {
return BindingBuilder
.bind(TestDirectQueue())
.to(TestDirectExchange())
.with(RabbitMqConstant.TEST_DIRECT_ROUTING); // 精确匹配
}
}
消费者 (DirectReceiver.java)
@Component
public class DirectReceiver {
Integer count = 0;
@RabbitHandler
public void process(Map testMessage) {
count++;
System.out.println("【Direct】消费者收到消息 [" + count + "]: " + testMessage.toString());
// 注意:若全局配置了 manual ack,此处需配合 MessageListenerConfig 使用
}
}
3.2 主题交换机 (Topic Exchange)
- 机制: 模糊匹配。支持通配符:
*(星号): 匹配一个单词 (如user.*匹配user.login)。#(井号): 匹配零个或多个单词 (如user.#匹配user.login.success)。
- 场景: 复杂的日志收集、多级业务通知。
测试类 (TopicSender.java)
@SpringBootTest
class TopicTests {
@Autowired
RabbitTemplate rabbitTemplate;
/**
* queue: topic.man 和 topic.woman 都可以发送消息
* exchange: topicExchange
* bindingKey: topic.man 和 topic.# 可以接收到消息
*
* 创建两个 JAVA 的接收器,一个接收 topic.man,一个接收 topic.woman
* JAVA 接收器 topic.man 可以绑定,
* topic.#
* topic.man
* JAVA 接收器 topic.woman 只能绑定
* topic.#
*/
@Test
public void sendTopicMessage1() {
sendMan();
}
@Test
public void sendTopicMessage2() {
sendWoMan();
}
@Test
public void sendTopicMessage() throws InterruptedException {
while (true) {
// sendMan();
sendWoMan();
sleep(100);
}
}
private void sendMan() {
String messageId = String.valueOf(UUID.randomUUID());
String messageData = "message: M A N ";
String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
Map<String, Object> manMap = new HashMap<>();
manMap.put("messageId", messageId);
manMap.put("messageData", messageData);
manMap.put("createTime", createTime);
rabbitTemplate.convertAndSend("topicExchange", "topic.man", manMap);
}
private void sendWoMan() {
String messageId = String.valueOf(UUID.randomUUID());
String messageData = "message: woman is all ";
String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
Map<String, Object> womanMap = new HashMap<>();
womanMap.put("messageId", messageId);
womanMap.put("messageData", messageData);
womanMap.put("createTime", createTime);
rabbitTemplate.convertAndSend("topicExchange", "topic.woman", womanMap);
}
}
配置类 (TopicRabbitConfig.java)
@Configuration
public class TopicRabbitConfig {
@Bean
public Queue firstQueue() {
return new Queue(RabbitMqConstant.TOPIC_ROUTING_KEY_MAN); // topic.man
}
@Bean
public Queue secondQueue() {
return new Queue(RabbitMqConstant.TOPIC_ROUTING_KEY_WOMAN); // topic.woman
}
@Bean
TopicExchange exchange() {
return new TopicExchange(RabbitMqConstant.TOPIC_EXCHANGE);
}
// 绑定:只接收 topic.man
@Bean
Binding bindingExchangeMessageByMan() {
return BindingBuilder.bind(firstQueue()).to(exchange()).with(RabbitMqConstant.TOPIC_ROUTING_KEY_MAN);
}
// 绑定:接收 topic.# (所有以 topic. 开头的)
@Bean
Binding bindingExchangeMessageAll() {
return BindingBuilder.bind(secondQueue()).to(exchange()).with(RabbitMqConstant.TOPIC_ROUTING_KEY_ALL);
}
}
消费者 (TopicReceiver.java)
@Component
@RabbitListener(queues = "topic.man")
public class TopicManReceiver {
@RabbitHandler
public void process(Map testMessage) {
System.out.println("【man】TopicManReceiver消费者收到消息 : " + testMessage.toString());
}
}
@Component
@RabbitListener(queues = "topic.woman")
public class TopicTotalReceiver {
@RabbitHandler
public void process(Map testMessage) {
System.out.println("【#】TopicTotalReceiver消费者收到消息 : " + testMessage.toString());
}
}
3.3 扇形交换机 (Fanout Exchange)
- 机制: 广播。忽略 Routing Key,发送给所有绑定的队列。
- 场景: 系统广播、缓存失效通知。
- 特点: 速度最快,但灵活性最低。
测试类 (FanoutSender.java)
@SpringBootTest
class FanoutTests {
@Autowired
RabbitTemplate rabbitTemplate; // 使用RabbitTemplate,这提供了接收/发送等等方法
@Test
public void sendFanoutMessage() throws InterruptedException {
while (true) {
String messageId = String.valueOf(UUID.randomUUID());
String messageData = "message: testFanoutMessage ";
String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
Map<String, Object> map = new HashMap<>();
map.put("messageId", messageId);
map.put("messageData", messageData);
map.put("createTime", createTime);
rabbitTemplate.convertAndSend("fanoutExchange", null, map);
sleep(100);
}
}
}
配置类 (FanoutRabbitConfig.java)
@Configuration
public class FanoutRabbitConfig {
@Bean public Queue queueA() { return new Queue("fanout.A"); }
@Bean public Queue queueB() { return new Queue("fanout.B"); }
@Bean public Queue queueC() { return new Queue("fanout.C"); }
@Bean
FanoutExchange fanoutExchange() {
return new FanoutExchange("fanoutExchange");
}
@Bean Binding bindingExchangeA() { return BindingBuilder.bind(queueA()).to(fanoutExchange()); }
@Bean Binding bindingExchangeB() { return BindingBuilder.bind(queueB()).to(fanoutExchange()); }
@Bean Binding bindingExchangeC() { return BindingBuilder.bind(queueC()).to(fanoutExchange()); }
}
消费者 (FanoutReceiver.java)
@Component
@RabbitListener(queues = "fanout.A")
public class FanoutReceiverA {
@RabbitHandler
public void process(Map testMessage) {
System.out.println("FanoutReceiverA消费者收到消息 : " +testMessage.toString());
}
}
@Component
@RabbitListener(queues = "fanout.B")
public class FanoutReceiverB {
@RabbitHandler
public void process(Map testMessage) {
System.out.println("FanoutReceiverB消费者收到消息 : " +testMessage.toString());
}
}
@Component
@RabbitListener(queues = "fanout.C")
public class FanoutReceiverC {
@RabbitHandler
public void process(Map testMessage) {
System.out.println("FanoutReceiverC消费者收到消息 : " + testMessage.toString());
}
}
第四部分:高级特性 - 死信与延迟队列
4.1 死信队列 (Dead Letter Exchange, DLX)
原理: 当消息在普通队列中变成“死信”时,自动转发到指定的死信交换机。
触发条件:
- 消息被拒绝 (
basic.reject/basic.nack) 且requeue=false。 - 消息在队列中过期 (TTL)。
- 队列达到最大长度,消息被挤出。
测试类 (DeadLetterSender.java)
@SpringBootTest
public class DelayTests {
@Autowired
RabbitTemplate rabbitTemplate; // 使用RabbitTemplate,这提供了接收/发送等等方法
@Test
public void testDelay() throws InterruptedException {
while (true) {
String string = UUID.randomUUID().toString();
String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
Map<String, Object> map = new HashMap<>();
map.put("messageId",string);
map.put("messageData","订单信息:id=1,time=2022年03月30日11:41:47");
map.put("createTime", createTime);
CorrelationData correlationData = new CorrelationData();
rabbitTemplate.convertAndSend(DelayRabbitConfig.EXCHANGE_NAME,
"test.order.msg", map.toString().getBytes(StandardCharsets.UTF_8),correlationData);
Thread.sleep(100);
}
}
}
配置类 (DeadLetterRabbitConfig.java)
@Configuration
public class DeadLetterRabbitConfig {
public static final String EXCHANGE_NAME_DLX = "exchange_dlx6";
public static final String QUEUE_NAME_DLX = "queue_dlx6";
public static final String EXCHANGE_NAME = "test_exchange_dlx6";
public static final String QUEUE_NAME = "test_queue_dlx6";
// 1. 普通交换机
@Bean("test_exchange_dlx")
public Exchange bootExchange(){
return ExchangeBuilder.topicExchange(EXCHANGE_NAME).durable(true).build();
}
// 2. 死信交换机
@Bean("exchange_dlx")
public Exchange dlxExchange(){
return ExchangeBuilder.topicExchange(EXCHANGE_NAME_DLX).durable(true).build();
}
// 3. 普通队列 (关键:设置死信参数)
@Bean("test_queue_dlx")
public Queue bootQueue(){
Map<String, Object> args = new HashMap<>();
args.put("x-dead-letter-exchange", EXCHANGE_NAME_DLX); // 死信交换机
args.put("x-dead-letter-routing-key", "dlx.#"); // 死信路由键
args.put("x-message-ttl", 10000); // 消息 TTL 10 秒
args.put("x-max-length", 10); // 队列最大长度 10
return QueueBuilder.durable(QUEUE_NAME).withArguments(args).build();
}
// 4. 死信队列
@Bean("queue_dlx")
public Queue dlxQueue(){
return QueueBuilder.durable(QUEUE_NAME_DLX).build();
}
// 5. 绑定关系
@Bean
public Binding bindDlxQueue(@Qualifier("queue_dlx") Queue queue,
@Qualifier("exchange_dlx") Exchange exchange){
return BindingBuilder.bind(queue).to(exchange).with("dlx.#").noargs();
}
@Bean
public Binding bindNormalQueue(@Qualifier("test_queue_dlx") Queue queue,
@Qualifier("test_exchange_dlx") Exchange exchange){
return BindingBuilder.bind(queue).to(exchange).with("test.dlx.#").noargs();
}
}
4.2 延迟队列 (Delay Queue) - 双方案详解
方案一:TTL + 死信队列 (通用方案,无需插件)
原理: 利用消息过期成为死信的机制。
缺点: 如果队列头部消息未过期,会阻塞后面已过期的消息(FIFO 限制)。
配置: 参考 DelayRabbitConfig.java (逻辑同死信配置,重点设置 x-message-ttl)。
消费者: 必须监听死信队列。
@Component
public class DelayReceiver implements ChannelAwareMessageListener {
@Override
public void onMessage(Message message, Channel channel) throws Exception {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
String body = new String(message.getBody(), StandardCharsets.UTF_8);
System.out.println("⏰ [" + LocalDateTime.now() + "] 接收到延迟消息: " + body);
// 业务逻辑:检查订单状态,若未支付则取消
// ...
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
// 拒绝且不重回队列,防止死循环
channel.basicNack(deliveryTag, false, false);
e.printStackTrace();
}
}
}
方案二:RabbitMQ 延迟插件 (推荐,高性能)
原理: 安装 rabbitmq_delayed_message_exchange 插件,交换机类型变为 x-delayed-message。
优点: 消息在交换机层面延迟,不阻塞队列,精度更高。
步骤:
- 下载对应版本的
.ez插件放入plugins目录。 - 启用插件:
rabbitmq-plugins enable rabbitmq_delayed_message_exchange。 - 代码修改:
@Bean
public CustomExchange delayExchange() {
Map<String, Object> args = new HashMap<>();
args.put("x-delayed-type", "direct"); // 内部使用 direct 交换
return new CustomExchange("delayed-exchange", "x-delayed-message", true, false, args);
}
// 发送时指定延迟时间
Map<String, Object> headers = new HashMap<>();
headers.put("x-delay", 5000); // 延迟 5000ms
AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder().headers(headers).build();
rabbitTemplate.send("delayed-exchange", "routing-key", new Message("msg".getBytes(), properties));
第五部分:消息可靠性与手动确认
5.1 为什么需要手动确认?
默认 AUTO 模式下,消费者一旦收到消息,RabbitMQ 就认为已处理并删除消息。若业务逻辑报错,消息将永久丢失。
Manual 模式: 只有代码显式调用 basicAck,消息才会被删除。
5.2 手动确认配置 (MessageListenerConfig.java)
@Configuration
public class MessageListenerConfig {
@Autowired
private CachingConnectionFactory connectionFactory;
@Autowired
private MyAckReceiver myAckReceiver;
@Bean
public SimpleMessageListenerContainer container() {
SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory);
container.setConcurrentConsumers(1);
container.setMaxConcurrentConsumers(5); // 动态扩容
container.setAcknowledgeMode(AcknowledgeMode.MANUAL); // 核心:手动确认
container.setQueueNames("TestDirectQueue", "fanout.A", "order_que_dlx_ttl");
container.setMessageListener(myAckReceiver);
return container;
}
}
5.3 统一接收者 (MyAckReceiver.java)
@Component
public class MyAckReceiver implements ChannelAwareMessageListener {
@Override
public void onMessage(Message message, Channel channel) throws Exception {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
String queueName = message.getMessageProperties().getConsumerQueue();
try {
String body = new String(message.getBody(), StandardCharsets.UTF_8);
System.out.println("📥 收到消息 [队列:" + queueName + "]: " + body);
// === 业务逻辑 ===
if (body.contains("error")) {
throw new RuntimeException("模拟业务异常");
}
// === 成功:手动 ACK ===
channel.basicAck(deliveryTag, false);
System.out.println("✅ 消息确认成功");
} catch (Exception e) {
e.printStackTrace();
// === 失败策略 ===
// 方案 B: 拒绝并不再入队 (requeue=false) -> 进入死信队列 (推荐)
channel.basicNack(deliveryTag, false, false);
System.out.println("❌ 消息处理失败,已送入死信队列");
}
}
}
第六部分:控制器测试与交互流程
6.1 测试控制器 (RabbitMQController.java)
@RestController
@RequestMapping("/mq")
public class RabbitMQController {
@Autowired
private RabbitTemplate rabbitTemplate;
// 1. 发送直连消息
@GetMapping("/sendDirect")
public String sendDirect() {
Map<String, Object> msg = buildMsg("Hello Direct");
rabbitTemplate.convertAndSend(RabbitMqConstant.TEST_DIRECT_EXCHANGE,
RabbitMqConstant.TEST_DIRECT_ROUTING, msg);
return "Direct 消息已发送";
}
// 2. 发送主题消息
@GetMapping("/sendTopic/{type}")
public String sendTopic(@PathVariable String type) {
Map<String, Object> msg = buildMsg("Hello Topic: " + type);
String routingKey = "topic." + type;
rabbitTemplate.convertAndSend(RabbitMqConstant.TOPIC_EXCHANGE, routingKey, msg);
return "Topic 消息已发送: " + routingKey;
}
// 3. 发送广播消息
@GetMapping("/sendFanout")
public String sendFanout() {
Map<String, Object> msg = buildMsg("Hello Fanout Broadcast");
rabbitTemplate.convertAndSend("fanoutExchange", "", msg);
return "Fanout 广播已发送";
}
// 4. 发送延迟消息 (触发死信)
@GetMapping("/sendDelay")
public String sendDelay() {
Map<String, Object> msg = buildMsg("Order ID: 1001");
rabbitTemplate.convertAndSend("test_exchange_dlx_ttl", "test.order.delay", msg);
return "延迟消息已发送,预计 10 秒后处理";
}
private Map<String, Object> buildMsg(String data) {
Map<String, Object> map = new HashMap<>();
map.put("messageId", UUID.randomUUID().toString());
map.put("messageData", data);
map.put("createTime", LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")));
return map;
}
}
6.2 延迟队列交互流程图
第七部分:生产环境最佳实践与避坑
7.1 消息可靠性保障 (三板斧)
- 生产者: 开启
ConfirmCallback+ReturnCallback,失败时记录 DB 并重试。 - Broker: 队列、交换机、消息均设置为持久化 (
durable=true,deliveryMode=2)。 - 消费者: 使用
Manual Ack,业务成功后再 ACK。
7.2 消息重复消费 (幂等性)
网络抖动可能导致 ACK 丢失,RabbitMQ 会重发消息。
- 解决方案: 业务层必须实现幂等性。
- 手段:
- 数据库唯一索引: 插入前检查
message_id是否存在。 - Redis 原子操作:
setnx key value,处理前先占位。 - 状态机判断: 订单状态只能从
PAYING->PAID,不能重复变。
- 数据库唯一索引: 插入前检查
7.3 消息积压处理
- 临时方案: 编写临时消费者,只负责把消息取出来转发到新的 Topic,然后启动大量新消费者并行处理。
- 根本方案: 优化消费者逻辑,增加消费者实例数量 (
concurrentConsumers)。
7.4 常见异常与解决
MessageConversionException: 通常是生产者发了 JSON,消费者没配Jackson2JsonMessageConverter。ChannelClosedException: 消费者代码抛出异常未捕获,导致通道关闭。务必在onMessage中try-catch。NOT_FOUND: 队列不存在。检查@Bean是否被扫描。
更多推荐




所有评论(0)