SpringBoot集成RocketMQ事務(wù)消息的完整指南
事務(wù)消息是 RocketMQ 提供的一種高級(jí)消息類型,用于解決分布式場(chǎng)景下,本地?cái)?shù)據(jù)庫(kù)事務(wù)與消息發(fā)送之間的一致性問(wèn)題。它通過(guò)兩階段提交和事務(wù)狀態(tài)回查機(jī)制,確保本地事務(wù)執(zhí)行與消息投遞達(dá)到最終一致性,尤其適用于訂單支付、積分變更等需要高可靠性的業(yè)務(wù)場(chǎng)景。
事務(wù)消息的核心原理
事務(wù)消息的核心機(jī)制可以概括為以下兩個(gè)階段和一種補(bǔ)償機(jī)制:
第一階段:發(fā)送半消息(Half Message)??
- 生產(chǎn)者向 Broker 發(fā)送一條半消息。這條消息與普通消息不同,它已經(jīng)持久化到服務(wù)端,但對(duì)消費(fèi)者不可見,暫時(shí)不能被消費(fèi)。
- Broker 收到半消息并持久化成功后,會(huì)向生產(chǎn)者返回確認(rèn)響應(yīng)。
第二階段:提交或回滾?
生產(chǎn)者開始執(zhí)行本地事務(wù)?(例如,操作本地?cái)?shù)據(jù)庫(kù))。
根據(jù)本地事務(wù)的執(zhí)行結(jié)果(成功或失?。a(chǎn)者向 Broker 發(fā)送 ?二次確認(rèn)指令?(Commit 或 Rollback)。
- ?Commit?:Broker 將半消息轉(zhuǎn)換為正式消息,對(duì)消費(fèi)者可見,等待被消費(fèi)。
- ?Rollback?:Broker 會(huì)回滾該事務(wù),即刪除半消息,消息不會(huì)被投遞。
?事務(wù)回查(Transaction Check)??
- 這是關(guān)鍵的補(bǔ)償機(jī)制。如果因?yàn)榫W(wǎng)絡(luò)閃斷、生產(chǎn)者應(yīng)用重啟等原因,導(dǎo)致 Broker 長(zhǎng)時(shí)間未收到二次確認(rèn),Broker 會(huì)主動(dòng)向生產(chǎn)者發(fā)起消息回查。
- 生產(chǎn)者收到回查后,需要去檢查該消息對(duì)應(yīng)的本地事務(wù)的最終狀態(tài)(例如查詢數(shù)據(jù)庫(kù)),并根據(jù)實(shí)際狀態(tài)再次向 Broker 提交 Commit 或 Rollback 指令。這保證了即使在異常情況下,事務(wù)也能最終達(dá)成一致。
為了更直觀地理解整個(gè)流程,下圖概括了事務(wù)消息的完整生命周期:

在SpringBoot項(xiàng)目中實(shí)現(xiàn)事務(wù)消息
下面我們基于 rocketmq-spring-boot-starter來(lái)實(shí)現(xiàn)一個(gè)完整的事務(wù)消息示例,以“訂單支付成功后通知積分服務(wù)增加積分”為場(chǎng)景。
1. 添加依賴
首先確保 pom.xml中包含必要的依賴:
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.2.3</version>
</dependency>
2. 配置生產(chǎn)者與事務(wù)監(jiān)聽器
核心是創(chuàng)建一個(gè)事務(wù)監(jiān)聽器,它包含了執(zhí)行本地事務(wù)和處理事務(wù)回查的兩個(gè)方法。
import org.apache.rocketmq.spring.annotation.RocketMQTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionState;
import org.springframework.messaging.Message;
@Service
@RocketMQTransactionListener(txProducerGroup = "tx-order-group") // 與發(fā)送方組名一致
public class OrderTransactionListenerImpl implements RocketMQLocalTransactionListener {
@Autowired
private OrderService orderService;
/**
* 執(zhí)行本地事務(wù)
* @param msg 收到的消息
* @param arg 調(diào)用sendMessageInTransaction時(shí)傳入的額外參數(shù)
* @return 事務(wù)狀態(tài)
*/
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 從消息頭或arg中獲取業(yè)務(wù)ID,如訂單ID
String orderId = (String) msg.getHeaders().get("orderId");
try {
// 執(zhí)行本地業(yè)務(wù)邏輯,例如:更新訂單狀態(tài)為“支付成功”
boolean success = orderService.updateOrderStatus(orderId, OrderStatus.PAID);
// 根據(jù)執(zhí)行結(jié)果返回提交或回滾
return success ? RocketMQLocalTransactionState.COMMIT : RocketMQLocalTransactionState.ROLLBACK;
} catch (Exception e) {
// 記錄日志,返回UNKNOWN狀態(tài),等待Broker回查
return RocketMQLocalTransactionState.UNKNOWN;
}
}
/**
* 事務(wù)回查方法
* @param msg 收到的消息
* @return 事務(wù)狀態(tài)
*/
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
String orderId = (String) msg.getHeaders().get("orderId");
// 根據(jù)orderId查詢數(shù)據(jù)庫(kù),確認(rèn)本地事務(wù)的最終狀態(tài)
OrderStatus status = orderService.queryOrderStatus(orderId);
if (OrderStatus.PAID.equals(status)) {
// 本地事務(wù)已成功,提交消息
return RocketMQLocalTransactionState.COMMIT;
} else if (OrderStatus.FAILED.equals(status)) {
// 本地事務(wù)已失敗,回滾消息
return RocketMQLocalTransactionState.ROLLBACK;
} else {
// 狀態(tài)仍不明確,繼續(xù)等待下次回查
return RocketMQLocalTransactionState.UNKNOWN;
}
}
}
3. 發(fā)送事務(wù)消息
在業(yè)務(wù)服務(wù)中,使用 RocketMQTemplate發(fā)送事務(wù)消息。
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service;
@Service
public class OrderService {
@Autowired
private RocketMQTemplate rocketMQTemplate;
public void payOrder(String orderId) {
// 1. 構(gòu)建消息
Message<String> message = MessageBuilder.withPayload("訂單支付成功,增加積分")
.setHeader("orderId", orderId) // 設(shè)置業(yè)務(wù)ID,用于回查
.build();
// 2. 發(fā)送事務(wù)消息
// 參數(shù)1: 事務(wù)組名(需與監(jiān)聽器內(nèi)txProducerGroup一致)
// 參數(shù)2: 主題(Topic)
// 參數(shù)3: 消息體
// 參數(shù)4: 可選參數(shù),會(huì)傳遞給executeLocalTransaction方法的arg參數(shù)
TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction("tx-order-group", "order-topic", message, orderId);
System.out.println("發(fā)送結(jié)果:" + result.getSendStatus());
}
}
4. 消費(fèi)者端實(shí)現(xiàn)冪等性
事務(wù)消息只能保證消息生產(chǎn)端的一致性,消費(fèi)端需要自行保證消息的冪等性,因?yàn)榫W(wǎng)絡(luò)重試可能導(dǎo)致消息被重復(fù)消費(fèi)。
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Service;
@Service
@RocketMQMessageListener(topic = "order-topic", consumerGroup = "order-consumer-group")
public class OrderConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
// 1. 解析消息,獲取訂單ID
// 2. 【關(guān)鍵】?jī)绲刃孕r?yàn):查詢數(shù)據(jù)庫(kù)或Redis,判斷該訂單的積分是否已經(jīng)添加過(guò)
// if (已處理) { return; }
// 3. 執(zhí)行業(yè)務(wù)邏輯(例如,為用戶增加積分)
// creditService.addCredit(...);
// 4. 記錄處理狀態(tài),標(biāo)記該消息已處理
}
}
重要注意事項(xiàng)與最佳實(shí)踐
- ?事務(wù)狀態(tài)的可查詢性?:
checkLocalTransaction方法需要能夠查詢本地事務(wù)的最終狀態(tài)。通常的做法是,在執(zhí)行本地事務(wù)時(shí),將事務(wù)狀態(tài)(如訂單狀態(tài))持久化到數(shù)據(jù)庫(kù)中,以便回查時(shí)使用。 - ?避免未知狀態(tài)?:雖然
UNKNOWN狀態(tài)是回查機(jī)制的保障,但在生產(chǎn)中應(yīng)盡量明確返回COMMIT或ROLLBACK,避免大量消息進(jìn)入回查流程,影響系統(tǒng)性能和增加復(fù)雜度。 - ?消息回查配置?:Broker 端有關(guān)于回查間隔和最大回查次數(shù)的配置,需要根據(jù)業(yè)務(wù)容忍度進(jìn)行合理設(shè)置。
- ?主題類型匹配?:事務(wù)消息必須發(fā)送到類型為
Transaction的主題上。
以上就是SpringBoot集成RocketMQ事務(wù)消息的完整指南的詳細(xì)內(nèi)容,更多關(guān)于SpringBoot集成RocketMQ事務(wù)消息的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!
相關(guān)文章
SpringCloud解決feign調(diào)用token丟失問(wèn)題解決辦法
在feign調(diào)用中可能會(huì)遇到如下問(wèn)題:同步調(diào)用中,token丟失,這種可以通過(guò)創(chuàng)建一個(gè)攔截器,將token做透?jìng)鱽?lái)解決,異步調(diào)用中,token丟失,這種就無(wú)法直接透?jìng)髁?因?yàn)樽泳€程并沒(méi)有token,這種需要先將token從父線程傳遞到子線程,再進(jìn)行透?jìng)?/div> 2024-05-05
Java獲取當(dāng)前操作系統(tǒng)的信息實(shí)例代碼
這篇文章主要介紹了Java獲取當(dāng)前操作系統(tǒng)的信息實(shí)例代碼,具有一定借鑒價(jià)值,需要的朋友可以參考下。2017-12-12
Knife4j?3.0.3?整合SpringBoot?2.6.4的詳細(xì)過(guò)程
本文要講的是?Knife4j?3.0.3?整合SpringBoot?2.6.4,在SpringBoot?2.4以上的版本和之前的版本還是不一樣的,這個(gè)也容易導(dǎo)致一些問(wèn)題,本文就這兩個(gè)版本的整合做一個(gè)實(shí)戰(zhàn)介紹2022-09-09
SpringBoot中實(shí)時(shí)監(jiān)控Redis命令流的實(shí)現(xiàn)
在Redis的日常使用和調(diào)試中,監(jiān)控命令流有助于我們更好地理解 Redis的工作狀態(tài),Redis提供了MONITOR命令,可以實(shí)時(shí)輸出Redis中所有客戶端的命令請(qǐng)求,本文將介紹如何使用Jedis實(shí)現(xiàn)這一功能,并對(duì)比telnet實(shí)現(xiàn)MONITOR機(jī)制的工作方式,需要的朋友可以參考下2024-11-11
基于SpringBoot+Mybatis實(shí)現(xiàn)Mysql分表
這篇文章主要為大家詳細(xì)介紹了基于SpringBoot+Mybatis實(shí)現(xiàn)Mysql分表的相關(guān)知識(shí),文中的示例代碼講解詳細(xì),感興趣的小伙伴可以跟隨小編一起學(xué)習(xí)一下2025-04-04最新評(píng)論

