SpringBoot+RabbitMQ實現(xiàn)消息可靠投遞+防重復(fù)消費(可直接落地)
提示:文章寫完后,目錄可以自動生成,如何生成可參考右邊的幫助文檔
前言
Spring Boot + RabbitMQ 實戰(zhàn):消息可靠投遞+防重復(fù)消費(可直接落地)
在高并發(fā)業(yè)務(wù)場景中,RabbitMQ 作為消息中間件,核心作用是削峰填谷、解耦服務(wù),但最關(guān)鍵的兩個問題的是:消息不丟失、不重復(fù)消費。
基于 Spring Boot 整合 RabbitMQ,提供一套可直接復(fù)制、生產(chǎn)環(huán)境可用的實戰(zhàn)代碼,涵蓋「生產(chǎn)者 Confirm 確認、消息持久化、消費者手動 ACK、Redis 冪等防重」全流程,避開所有常見坑,新手也能快速落地。
適用場景:訂單異步創(chuàng)建、短信/通知推送、物流狀態(tài)同步等所有需要保證消息可靠性的業(yè)務(wù),尤其適配高并發(fā)下單、秒殺等場景。
一、核心需求與技術(shù)選型
1. 核心需求(必滿足)
- 生產(chǎn)者:消息必須送達 RabbitMQ Broker,失敗可重試,杜絕生產(chǎn)者丟消息
- 消息本身:Broker 重啟、服務(wù)宕機后,消息不丟失
- 消費者:業(yè)務(wù)處理成功后再確認消息,異常可重新入隊,杜絕消費端丟消息
- 冪等性:避免因網(wǎng)絡(luò)重試、消息重入隊導(dǎo)致的重復(fù)消費(比如重復(fù)創(chuàng)建訂單、重復(fù)扣庫存)
2. 技術(shù)選型
- 框架:Spring Boot 2.x(兼容 3.x,只需微調(diào)依賴)
- 消息中間件:RabbitMQ 3.9+
- 冪等校驗:Redis(高效判重)+ 數(shù)據(jù)庫唯一索引(兜底)
- 核心依賴:spring-boot-starter-amqp、spring-boot-starter-data-redis
二、環(huán)境配置(application.yml)
核心配置:開啟生產(chǎn)者 Confirm 機制、Return 機制,消費者手動 ACK,同時配置限流防止數(shù)據(jù)庫被沖垮,注釋清晰可直接復(fù)制。
spring:
# RabbitMQ 核心配置
rabbitmq:
host: 127.0.0.1 # 本地環(huán)境,生產(chǎn)環(huán)境替換為服務(wù)器地址
port: 5672 # RabbitMQ 默認端口
username: guest # 默認用戶名,生產(chǎn)環(huán)境需修改為自定義賬號
password: guest # 默認密碼,生產(chǎn)環(huán)境需修改
virtual-host: / # 虛擬主機,默認即可
connection-timeout: 10000 # 連接超時時間,避免無限等待
# 1. 生產(chǎn)者確認機制:確保消息到達 Broker
publisher-confirm-type: correlated # correlated:異步回調(diào),獲取確認結(jié)果
# 2. 消息回退機制:消息無法路由時返回生產(chǎn)者,避免消息丟失
publisher-returns: true
# 3. 消費者配置
listener:
simple:
acknowledge-mode: manual # 手動 ACK(關(guān)鍵!避免自動確認丟消息)
concurrency: 5 # 消費者核心并發(fā)數(shù)
max-concurrency: 10 # 消費者最大并發(fā)數(shù)
prefetch: 10 # 限流:每次只獲取10條消息,防止消費過快沖垮數(shù)據(jù)庫三、核心代碼實現(xiàn)(全可復(fù)制)
1. 依賴導(dǎo)入(pom.xml)
無需額外配置,導(dǎo)入 Spring Boot 整合 RabbitMQ 和 Redis 的 starter 即可。
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<!-- lombok 簡化代碼,可選但推薦 -->
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>2. 實體類(OrderCreateDTO)
訂單消息傳輸?shù)膶嶓w類,根據(jù)自身業(yè)務(wù)調(diào)整字段,需實現(xiàn) Serializable 接口(RabbitMQ 消息傳輸要求)。
import lombok.Data;
import java.io.Serializable;
/**
* 訂單創(chuàng)建消息DTO
*/
@Data
public class OrderCreateDTO implements Serializable {
// 訂單唯一編號(用于業(yè)務(wù)冪等)
private String orderSn;
// 用戶ID
private Long userId;
// 訂單金額
private BigDecimal orderAmount;
// 商品ID(多個可改為List)
private Long productId;
// 購買數(shù)量
private Integer quantity;
}3. 生產(chǎn)者:消息可靠投遞(Confirm + 持久化)
核心邏輯:
- 生成全局唯一 msgId(用于后續(xù)冪等判重)
- 設(shè)置消息持久化(Broker 重啟后消息不丟失)
- 開啟 Confirm 回調(diào),監(jiān)聽消息是否成功送達 Broker,失敗可重試/落庫
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.MessageDeliveryMode;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
import java.util.UUID;
/**
* 訂單消息生產(chǎn)者(可靠投遞)
*/
@Component
@Slf4j
@RequiredArgsConstructor
public class OrderProducer {
// 注入RabbitTemplate,用于發(fā)送消息
private final RabbitTemplate rabbitTemplate;
// 交換機名稱(需與消費者隊列綁定)
private static final String ORDER_EXCHANGE = "order.exchange";
// 路由鍵(需與隊列綁定,確保消息能路由到指定隊列)
private static final String ORDER_CREATE_ROUTING_KEY = "order.create";
/**
* 發(fā)送訂單創(chuàng)建消息
* @param dto 訂單創(chuàng)建DTO
*/
public void sendOrderMsg(OrderCreateDTO dto) {
// 1. 生成全局唯一消息ID,用于冪等判重(UUID保證唯一性)
String msgId = UUID.randomUUID().toString().replace("-", "");
// 2. 關(guān)聯(lián)消息ID,用于Confirm回調(diào)獲取消息標識
CorrelationData correlationData = new CorrelationData(msgId);
// 3. 發(fā)送消息(設(shè)置持久化 + 攜帶消息ID)
rabbitTemplate.convertAndSend(
ORDER_EXCHANGE, // 交換機
ORDER_CREATE_ROUTING_KEY,// 路由鍵
dto, // 消息內(nèi)容
message -> {
// 設(shè)置消息持久化(MessageDeliveryMode.PERSISTENT:持久化)
message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
// 將msgId存入消息屬性,供消費者獲取
message.getMessageProperties().setMessageId(msgId);
return message;
},
correlationData // 關(guān)聯(lián)消息ID,用于Confirm回調(diào)
);
// 4. 生產(chǎn)者Confirm回調(diào):監(jiān)聽消息是否成功送達Broker
rabbitTemplate.setConfirmCallback((correlation, ack, cause) -> {
// 獲取回調(diào)的消息ID
String msgIdCallback = correlation.getId();
if (ack) {
// ack為true:消息成功送達Broker
log.info("消息發(fā)送成功,msgId:{}", msgIdCallback);
} else {
// ack為false:消息發(fā)送失敗
log.error("消息發(fā)送失敗,msgId:{},失敗原因:{}", msgIdCallback, cause);
// 失敗處理:可重試發(fā)送(建議最多3次),或入庫定時重發(fā)(避免消息丟失)
retrySendMsg(dto, msgIdCallback);
}
});
}
/**
* 消息發(fā)送失敗重試(簡單重試邏輯,可根據(jù)業(yè)務(wù)優(yōu)化)
*/
private void retrySendMsg(OrderCreateDTO dto, String msgId) {
int retryCount = 3; // 重試3次
for (int i = 0; i < retryCount; i++) {
try {
Thread.sleep(1000 * (i + 1)); // 指數(shù)退避重試(1s、2s、3s)
CorrelationData correlationData = new CorrelationData(msgId);
rabbitTemplate.convertAndSend(
ORDER_EXCHANGE,
ORDER_CREATE_ROUTING_KEY,
dto,
message -> {
message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
message.getMessageProperties().setMessageId(msgId);
return message;
},
correlationData
);
log.info("消息重試發(fā)送成功,msgId:{},重試次數(shù):{}", msgId, i + 1);
return;
} catch (Exception e) {
log.error("消息重試發(fā)送失敗,msgId:{},重試次數(shù):{}", msgId, i + 1, e);
if (i == retryCount - 1) {
// 重試3次仍失敗,入庫定時重發(fā)(此處省略入庫邏輯,可結(jié)合定時任務(wù)實現(xiàn))
log.error("消息重試3次失敗,msgId:{},已入庫待定時重發(fā)", msgId);
}
}
}
}
}4. 消費者:手動 ACK + Redis 冪等防重
核心邏輯:
- 手動 ACK:業(yè)務(wù)處理成功后,調(diào)用 basicAck 確認消息;異常則調(diào)用 basicNack 重新入隊
- Redis 冪等:用 setIfAbsent 存儲 msgId,已消費則直接 ACK,避免重復(fù)消費
- 業(yè)務(wù)兜底:結(jié)合數(shù)據(jù)庫唯一索引,防止 Redis 掛了導(dǎo)致的重復(fù)消費
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import com.rabbitmq.client.Channel;
import java.util.concurrent.TimeUnit;
/**
* 訂單消息消費者(手動ACK + 冪等防重)
*/
@Component
@Slf4j
@RequiredArgsConstructor
public class OrderConsumer {
// 注入Redis模板,用于冪等判重
private final StringRedisTemplate redisTemplate;
// 注入訂單服務(wù),處理核心業(yè)務(wù)邏輯
private final OrderService orderService;
// 隊列名稱(需與交換機、路由鍵綁定)
private static final String ORDER_CREATE_QUEUE = "order.create.queue";
// Redis 冪等鍵前綴(區(qū)分不同業(yè)務(wù)的消息)
private static final String MQ_CONSUMED_KEY_PREFIX = "mq:consumed:order:";
/**
* 消費訂單創(chuàng)建消息
* @param dto 消息內(nèi)容(自動反序列化)
* @param message 消息對象,用于獲取msgId
* @param channel 信道對象,用于手動ACK/NACK
*/
@RabbitListener(queues = ORDER_CREATE_QUEUE) // 監(jiān)聽指定隊列
public void consumeOrderMsg(OrderCreateDTO dto, Message message, Channel channel) throws Exception {
// 1. 獲取消息ID(生產(chǎn)者存入的msgId)
String msgId = message.getMessageProperties().getMessageId();
// 2. 獲取消息投遞標簽(用于手動ACK/NACK)
long deliveryTag = message.getMessageProperties().getDeliveryTag();
// 3. Redis 冪等判重:setIfAbsent 原子操作,避免并發(fā)重復(fù)消費
String redisKey = MQ_CONSUMED_KEY_PREFIX + msgId;
// 存入Redis,有效期24小時(根據(jù)業(yè)務(wù)調(diào)整,確保消息消費完成后不會被重復(fù)判斷)
Boolean consumeFlag = redisTemplate.opsForValue().setIfAbsent(redisKey, "1", 24, TimeUnit.HOURS);
// 4. 已消費過:直接ACK,避免重復(fù)處理
if (consumeFlag == null || !consumeFlag) {
log.warn("消息已消費,無需重復(fù)處理,msgId:{}", msgId);
// 手動ACK:deliveryTag為當(dāng)前消息標簽,false表示不批量確認
channel.basicAck(deliveryTag, false);
return;
}
try {
// 5. 處理核心業(yè)務(wù)邏輯:創(chuàng)建訂單、扣減庫存等
orderService.createOrder(dto);
// 6. 業(yè)務(wù)處理成功:手動ACK,通知RabbitMQ刪除消息
channel.basicAck(deliveryTag, false);
log.info("消息消費成功,msgId:{},訂單號:{}", msgId, dto.getOrderSn());
} catch (Exception e) {
log.error("消息消費異常,msgId:{},訂單號:{}", msgId, dto.getOrderSn(), e);
// 7. 異常處理:根據(jù)異常類型決定是否重入隊
// 可重試異常(如網(wǎng)絡(luò)波動、數(shù)據(jù)庫臨時不可用):重入隊(third參數(shù)為true)
// 不可重試異常(如業(yè)務(wù)校驗失敗、參數(shù)錯誤):直接拒絕,不重入隊(third參數(shù)為false)
if (isRetryException(e)) {
log.info("消息消費異常(可重試),將重入隊,msgId:{}", msgId);
channel.basicNack(deliveryTag, false, true);
} else {
log.info("消息消費異常(不可重試),直接拒絕,msgId:{}", msgId);
// 不可重試異常:拒絕消息,不重入隊(可結(jié)合死信隊列處理)
channel.basicReject(deliveryTag, false);
}
}
}
/**
* 判斷是否為可重試異常(根據(jù)自身業(yè)務(wù)調(diào)整)
*/
private boolean isRetryException(Exception e) {
// 示例:網(wǎng)絡(luò)異常、數(shù)據(jù)庫異??芍卦嚕瑯I(yè)務(wù)異常不可重試
return e instanceof RuntimeException
&& (e.getMessage().contains("網(wǎng)絡(luò)") || e.getMessage().contains("數(shù)據(jù)庫"));
}
}5. 業(yè)務(wù)層:Redis 冪等 + 數(shù)據(jù)庫兜底
**核心:**即使 Redis 掛了,通過數(shù)據(jù)庫唯一索引(order_sn)兜底,確保不會重復(fù)創(chuàng)建訂單、重復(fù)扣庫存。
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
/**
* 訂單服務(wù)(核心業(yè)務(wù)邏輯)
*/
@Service
@Slf4j
@RequiredArgsConstructor
public class OrderService {
private final OrderMapper orderMapper;
private final ProductMapper productMapper;
/**
* 創(chuàng)建訂單(帶冪等校驗)
* @param dto 訂單創(chuàng)建DTO
*/
@Transactional(rollbackFor = Exception.class) // 事務(wù)管理,異?;貪L
public void createOrder(OrderCreateDTO dto) {
String orderSn = dto.getOrderSn();
// 1. 數(shù)據(jù)庫冪等兜底:查詢訂單是否已存在(訂單表order_sn字段建唯一索引)
Integer orderCount = orderMapper.countByOrderSn(orderSn);
if (orderCount > 0) {
log.warn("訂單已存在,無需重復(fù)創(chuàng)建,訂單號:{}", orderSn);
return;
}
// 2. 核心業(yè)務(wù)邏輯:創(chuàng)建訂單、扣減庫存(根據(jù)自身業(yè)務(wù)實現(xiàn))
// ① 扣減商品庫存(需加鎖,避免超賣,此處省略分布式鎖邏輯)
Product product = productMapper.selectById(dto.getProductId());
if (product == null || product.getStock() < dto.getQuantity()) {
throw new RuntimeException("商品不存在或庫存不足,訂單號:" + orderSn);
}
product.setStock(product.getStock() - dto.getQuantity());
productMapper.updateById(product);
// ② 插入訂單記錄
Order order = new Order();
order.setOrderSn(orderSn);
order.setUserId(dto.getUserId());
order.setOrderAmount(dto.getOrderAmount());
order.setProductId(dto.getProductId());
order.setQuantity(dto.getQuantity());
order.setStatus(0); // 0:待支付
orderMapper.insert(order);
log.info("訂單創(chuàng)建成功,訂單號:{}", orderSn);
}
}6. 交換機、隊列綁定(可選,兩種方式)
方式1:代碼綁定(推薦,部署時自動創(chuàng)建,無需手動操作)
import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.DirectExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* RabbitMQ 交換機、隊列綁定配置
*/
@Configuration
public class RabbitMQConfig {
// 交換機名稱(與生產(chǎn)者、消費者一致)
public static final String ORDER_EXCHANGE = "order.exchange";
// 隊列名稱(與消費者一致)
public static final String ORDER_CREATE_QUEUE = "order.create.queue";
// 路由鍵(與生產(chǎn)者一致)
public static final String ORDER_CREATE_ROUTING_KEY = "order.create";
// 1. 聲明交換機(direct類型,持久化)
@Bean
public DirectExchange orderExchange() {
// durable=true:交換機持久化,重啟后不丟失
return new DirectExchange(ORDER_EXCHANGE, true, false);
}
// 2. 聲明隊列(持久化)
@Bean
public Queue orderCreateQueue() {
// durable=true:隊列持久化;exclusive=false:不排他;autoDelete=false:不自動刪除
return new Queue(ORDER_CREATE_QUEUE, true, false, false);
}
// 3. 綁定交換機、隊列、路由鍵
@Bean
public Binding orderCreateBinding() {
return BindingBuilder.bind(orderCreateQueue())
.to(orderExchange())
.with(ORDER_CREATE_ROUTING_KEY);
}
}方式2:RabbitMQ 管理界面手動綁定(適合測試環(huán)境,生產(chǎn)環(huán)境推薦代碼綁定)
- 登錄 RabbitMQ 管理界面(默認地址:http://localhost:15672)
- 創(chuàng)建交換機:類型 direct,名稱 order.exchange,勾選 durable
- 創(chuàng)建隊列:名稱 order.create.queue,勾選 durable
- 綁定:交換機 → 隊列,路由鍵填寫 order.create
四、關(guān)鍵知識點與避坑點(實戰(zhàn)重點)
消息不丟失的3個關(guān)鍵
- 生產(chǎn)者:開啟 Confirm 機制,確保消息到達 Broker,失敗重試/落庫
- 消息:設(shè)置持久化(MessageDeliveryMode.PERSISTENT),交換機、隊列也需持久化
- 消費者:手動 ACK,業(yè)務(wù)成功后再確認,異常合理處理(重入隊/死信)
防重復(fù)消費的2層保障
- 第一層:Redis setIfAbsent 原子操作(高效判重,適合高并發(fā))
- 第二層:數(shù)據(jù)庫唯一索引(兜底,防止 Redis 掛了導(dǎo)致的重復(fù)消費)
常見坑及解決方案
- 坑1:消息自動 ACK → 解決方案:配置 acknowledge-mode: manual,手動 ACK
- 坑2:消息未持久化 → 解決方案:設(shè)置消息、交換機、隊列均為持久化(durable=true)
- 坑3:Redis 掛了導(dǎo)致重復(fù)消費 → 解決方案:數(shù)據(jù)庫唯一索引兜底
- 坑4:消費者并發(fā)過高沖垮數(shù)據(jù)庫 → 解決方案:配置 prefetch 限流,控制每次獲取的消息數(shù)
- 坑5:消息發(fā)送失敗后不重試 → 解決方案:實現(xiàn) Confirm 回調(diào),失敗后指數(shù)退避重試,重試失敗入庫定時重發(fā)
五、測試驗證(快速驗證可用性)
- 啟動 RabbitMQ 服務(wù)(本地可通過 Docker 快速部署)
- 啟動 Spring Boot 項目,自動創(chuàng)建交換機、隊列并綁定
- 編寫測試類,調(diào)用 OrderProducer 的 sendOrderMsg 方法發(fā)送消息
- 查看日志:消息發(fā)送成功 → 消費成功 → 訂單創(chuàng)建成功
- 測試異常場景:關(guān)閉數(shù)據(jù)庫,發(fā)送消息,查看是否重入隊;恢復(fù)數(shù)據(jù)庫后,查看是否正常消費
- 測試重復(fù)消費:手動將消息重新入隊,查看是否會重復(fù)創(chuàng)建訂單(應(yīng)提示“訂單已存在”)
六、總結(jié)
本文提供的代碼的是生產(chǎn)環(huán)境真實落地版本,涵蓋了 RabbitMQ 消息可靠投遞和防重復(fù)消費的全流程,無需修改核心邏輯,只需根據(jù)自身業(yè)務(wù)調(diào)整實體類和業(yè)務(wù)方法,即可快速集成到項目中。
核心思路:生產(chǎn)者靠 Confirm 保送達,消息靠持久化保存活,消費者靠手動 ACK 保消費,冪等靠 Redis+數(shù)據(jù)庫保唯一,四者結(jié)合,徹底解決 RabbitMQ 消息丟失和重復(fù)消費的痛點。
后續(xù)可優(yōu)化方向:消息重試機制(結(jié)合定時任務(wù))、死信隊列(處理不可重試異常消息)、分布式鎖(防止庫存超賣),可根據(jù)業(yè)務(wù)復(fù)雜度逐步迭代
到此這篇關(guān)于SpringBoot+RabbitMQ實現(xiàn)消息可靠投遞+防重復(fù)消費(可直接落地)的文章就介紹到這了,更多相關(guān)SpringBoot RabbitMQ 消息可靠投遞+防重復(fù)消費內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
jquery uploadify和apache Fileupload實現(xiàn)異步上傳文件示例
這篇文章主要介紹了jquery uploadify和apache Fileupload實現(xiàn)異步上傳文件示例,需要的朋友可以參考下2014-05-05
Java Stream中的Spliterator類概念及原理解析
Spliterator是Java 8引入的一個接口,位于java.util包中,它結(jié)合了迭代器(Iterator)的遍歷能力和分割器(Splitter)的分割能力,本文將詳細介紹Spliterator的概念、原理、作用、類中定義的關(guān)鍵方法,以及它在Stream API中的實際應(yīng)用,感興趣的朋友一起看看吧2024-08-08
Spring Boot中使用 Spring Security 構(gòu)建權(quán)限系統(tǒng)的示例代碼
本篇文章主要介紹了Spring Boot中使用 Spring Security 構(gòu)建權(quán)限系統(tǒng)的示例代碼,具有一定的參考價值,有興趣的可以了解一下2017-08-08
Spring注解@Profile實現(xiàn)開發(fā)環(huán)境/測試環(huán)境/生產(chǎn)環(huán)境的切換
在進行軟件開發(fā)過程中,一般會將項目分為開發(fā)環(huán)境,測試環(huán)境,生產(chǎn)環(huán)境。本文主要介紹了Spring如何通過注解@Profile實現(xiàn)開發(fā)環(huán)境、測試環(huán)境、生產(chǎn)環(huán)境的切換,需要的可以參考一下2023-04-04
SpringBoot下獲取resources目錄下文件的常用方法
本文詳細介紹了SpringBoot獲取resources目錄下文件的常用方法,包括使用this.getClass()方法、ClassPathResource獲取以及hutool工具類ResourceUtil獲取,感興趣的可以了解一下2024-10-10

