RocketMQ中多消息不同狀態(tài)回查的設(shè)計與優(yōu)化過程
一、事務(wù)狀態(tài)回查的觸發(fā)條件
當(dāng)出現(xiàn)以下情況時,Broker 會主動發(fā)起事務(wù)狀態(tài)回查:
- 超時未確認(rèn):Producer 發(fā)送半消息后,在指定時間(
transactionTimeOut,默認(rèn) 60 秒)內(nèi)未發(fā)送 Commit/Rollback 指令 - Broker 重啟:Broker 重啟后,會恢復(fù)未完成的事務(wù)消息并觸發(fā)回查
- 超過最大提交延遲:半消息在 Broker 中存儲時間超過
transactionTimeout
二、多消息狀態(tài)回查的核心挑戰(zhàn)
- 狀態(tài)區(qū)分難題:多個消息可能同時處于
COMMIT、ROLLBACK、UNKNOW等不同狀態(tài),需精準(zhǔn)識別 - 并發(fā)控制需求:大量消息回查可能引發(fā)并發(fā)沖突,需保證狀態(tài)更新的原子性
- 性能優(yōu)化壓力:批量回查時若處理不當(dāng),可能導(dǎo)致 Broker 或 Producer 負(fù)載過高
三、狀態(tài)標(biāo)識與分類管理方案
1. 消息唯一標(biāo)識設(shè)計
- 業(yè)務(wù)主鍵綁定:在消息體中攜帶業(yè)務(wù)唯一標(biāo)識(如訂單 ID、交易號)
- 擴展屬性標(biāo)記:通過
Message.putUserProperty("bizType", "order")標(biāo)記消息類型
示例代碼:
// 發(fā)送消息時綁定業(yè)務(wù)標(biāo)識
Message msg = new Message("Topic", "Tag", "order123".getBytes());
msg.putUserProperty("bizId", "order123");
msg.putUserProperty("bizType", "order");
sendResult = producer.sendMessageInTransaction(transactionListener, msg, null);
2. 狀態(tài)分類存儲策略
| 存儲介質(zhì) | 適用場景 | 實現(xiàn)方式 |
|---|---|---|
| 數(shù)據(jù)庫 | 高可靠性要求,需持久化追溯 | 建表存儲 (bizId, status, updateTime),通過索引加速查詢 |
| Redis | 高性能讀寫,短期狀態(tài)存儲 | 使用 Hash 結(jié)構(gòu)存儲 {bizId: status},設(shè)置合理過期時間 |
| 本地緩存 | 高頻訪問,熱數(shù)據(jù)加速 | 結(jié)合 Guava Cache 或 ConcurrentHashMap,定期持久化到數(shù)據(jù)庫 |
3.回查機制的配置參數(shù)
| 參數(shù)名 | 默認(rèn)值 | 說明 |
|---|---|---|
| transactionTimeOut | 60 秒 | 事務(wù)超時時間,超過此時間未確認(rèn)則觸發(fā)回查 |
| transactionCheckMax | 15 次 | 最大回查次數(shù),超過此次數(shù)后 Broker 將根據(jù)策略處理(默認(rèn)丟棄消息) |
| transactionCheckInterval | 10 秒 | 兩次回查的時間間隔 |
4.回查實現(xiàn)的關(guān)鍵要點
冪等性設(shè)計:
- 回查方法可能被多次調(diào)用(如網(wǎng)絡(luò)波動導(dǎo)致 Broker 重復(fù)發(fā)起)
- 查詢操作必須是冪等的,避免重復(fù)提交或回滾
狀態(tài)存儲要求:
- 本地事務(wù)執(zhí)行后,必須將狀態(tài)持久化存儲(如數(shù)據(jù)庫、Redis)
- 回查時直接讀取持久化狀態(tài),而非依賴內(nèi)存變量
合理處理 UNKNOW 狀態(tài):
- 當(dāng)無法確定事務(wù)狀態(tài)時(如業(yè)務(wù)系統(tǒng)暫時不可用),返回
UNKNOW - Broker 會在配置的時間間隔后(
transactionCheckInterval)再次回查
避免長時間阻塞:
- 回查方法應(yīng)快速返回結(jié)果,避免長時間等待外部資源(如遠程服務(wù)調(diào)用)
- 若外部依賴不可用,建議先返回
UNKNOW,后續(xù)通過異步補償機制處理
5. 狀態(tài)機設(shè)計示例

四、多狀態(tài)回查的代碼實現(xiàn)
1. 基于消息屬性的差異化處理
public class MultiStatusTransactionListener implements TransactionListener {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 1. 解析消息屬性
String bizId = msg.getUserProperty("bizId");
String bizType = msg.getUserProperty("bizType");
// 2. 根據(jù)業(yè)務(wù)類型執(zhí)行不同本地事務(wù)
if ("order".equals(bizType)) {
return orderService.processOrder(bizId);
} else if ("payment".equals(bizType)) {
return paymentService.processPayment(bizId);
}
return LocalTransactionState.ROLLBACK_MESSAGE;
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 1. 解析消息屬性
String bizId = msg.getUserProperty("bizId");
String bizType = msg.getUserProperty("bizType");
// 2. 根據(jù)業(yè)務(wù)類型查詢不同狀態(tài)
if ("order".equals(bizType)) {
return orderService.checkOrderStatus(bizId);
} else if ("payment".equals(bizType)) {
return paymentService.checkPaymentStatus(bizId);
}
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
2. 批量回查優(yōu)化(減少網(wǎng)絡(luò)開銷)
// 自定義回查處理器,支持批量處理
public class BatchCheckProcessor {
// 緩存待回查消息,按業(yè)務(wù)類型分組
private final Map<String, List<String>> pendingCheck = new ConcurrentHashMap<>();
// 注冊回查消息
public void registerMessage(String bizType, String bizId) {
pendingCheck.computeIfAbsent(bizType, k -> new ArrayList<>()).add(bizId);
// 達到批量閾值或超時后觸發(fā)批量查詢
if (pendingCheck.get(bizType).size() >= 100 || needBatchCheck()) {
batchCheckAndClear(bizType);
}
}
// 批量查詢與狀態(tài)更新
private void batchCheckAndClear(String bizType) {
List<String> bizIds = pendingCheck.remove(bizType);
if (bizIds == null || bizIds.isEmpty()) return;
// 根據(jù)業(yè)務(wù)類型調(diào)用不同批量查詢接口
if ("order".equals(bizType)) {
Map<String, OrderStatus> statusMap = orderService.batchQueryStatus(bizIds);
// 批量更新狀態(tài)并發(fā)送響應(yīng)
statusMap.forEach((id, status) -> {
sendCheckResponse(id, mapToTransactionState(status));
});
}
// 其他業(yè)務(wù)類型處理...
}
}
五、多狀態(tài)回查的優(yōu)化策略
1. 按業(yè)務(wù)類型分組回查
Broker 配置:通過 transactionCheckListener 接口實現(xiàn)按主題或標(biāo)簽分組回查
示例配置:
<!-- 在 broker 配置文件中設(shè)置不同主題的回查策略 -->
<transactionCheckListener>
<topicCheckConfig>
<topic>order_topic</topic>
<checkInterval>5000</checkInterval> <!-- 訂單消息5秒回查一次 -->
<maxCheckTimes>20</maxCheckTimes>
</topicCheckConfig>
<topicCheckConfig>
<topic>payment_topic</topic>
<checkInterval>10000</checkInterval> <!-- 支付消息10秒回查一次 -->
<maxCheckTimes>10</maxCheckTimes>
</topicCheckConfig>
</transactionCheckListener>
2. 并發(fā)控制與限流
線程池隔離:為不同業(yè)務(wù)類型分配獨立的回查線程池
// 初始化多業(yè)務(wù)線程池
private final Map<String, ExecutorService> threadPools = new HashMap<>();
threadPools.put("order", new ThreadPoolExecutor(
10, 20, 60, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(1000),
new ThreadFactoryBuilder().setNameFormat("order-check-%d").build()
));
threadPools.put("payment", ...); // 支付業(yè)務(wù)線程池
信號量限流:控制同一時間回查的消息數(shù)量
private final Map<String, Semaphore> semaphores = new HashMap<>();
semaphores.put("order", new Semaphore(50)); // 訂單業(yè)務(wù)最多50個并發(fā)回查
3. 冪等性與防重處理
回查標(biāo)記:在狀態(tài)表中增加 check_version 字段,每次回查版本號遞增
分布式鎖:使用 Redis 或 Zookeeper 實現(xiàn)回查操作的全局鎖
// 回查前獲取分布式鎖,避免重復(fù)處理
boolean locked = redisTemplate.tryLock("check_lock:" + bizId, 3000);
if (locked) {
try {
// 執(zhí)行回查邏輯
} finally {
redisTemplate.unlock("check_lock:" + bizId);
}
}
六、多狀態(tài)回查的監(jiān)控與告警
1. 關(guān)鍵監(jiān)控指標(biāo)
| 指標(biāo)名稱 | 監(jiān)控目的 | 閾值建議 |
|---|---|---|
| 回查成功率 | 衡量回查處理有效性 | ≥99% |
| 平均回查耗時 | 評估系統(tǒng)處理性能 | ≤200ms |
| 待回查消息堆積量 | 發(fā)現(xiàn)潛在積壓風(fēng)險 | <1000 條 |
| 不同狀態(tài)消息占比 | 分析系統(tǒng)健康度 | COMMIT/ROLLBACK 占比 > 95% |
2. 告警策略示例
- 連續(xù)回查失敗告警:同一消息回查失敗超過 3 次時觸發(fā)
- 堆積超時告警:待回查消息在 Broker 中滯留超過
transactionTimeOut * 2時告警 - 業(yè)務(wù)類型異常告警:某類業(yè)務(wù)回查成功率連續(xù) 5 分鐘 < 80% 時告警
七、典型場景實現(xiàn)案例
電商訂單 - 支付聯(lián)動場景
消息類型:
- 訂單消息(bizType=order):回查間隔 5 秒,最大回查 20 次
- 支付消息(bizType=payment):回查間隔 10 秒,最大回查 10 次
狀態(tài)協(xié)同處理:
// 訂單狀態(tài)回查邏輯
public LocalTransactionState checkOrderStatus(String orderId) {
OrderStatus status = orderDao.getStatus(orderId);
if (status == SUCCESS) {
// 訂單成功時,主動檢查關(guān)聯(lián)的支付狀態(tài)
PaymentStatus payStatus = paymentDao.getStatusByOrder(orderId);
if (payStatus == SUCCESS) {
return COMMIT_MESSAGE;
} else {
// 支付未完成,延遲回查
return UNKNOW;
}
}
return mapToTransactionState(status);
}
最終一致性保障:
- 訂單狀態(tài)回查時,若發(fā)現(xiàn)支付未完成,觸發(fā)支付異步補償
- 支付狀態(tài)回查時,主動關(guān)聯(lián)訂單狀態(tài),確保兩者一致
通過以上方案,可有效處理多個消息的不同狀態(tài)回查,在保證最終一致性的同時,提升系統(tǒng)處理性能和穩(wěn)定性。實際應(yīng)用中需根據(jù)業(yè)務(wù)特性調(diào)整參數(shù)配置,并通過監(jiān)控持續(xù)優(yōu)化回查策略。
八、總結(jié)
以上為個人經(jīng)驗,希望能給大家一個參考,也希望大家多多支持腳本之家。
相關(guān)文章
Java使用POI-TL和JFreeChart動態(tài)生成Word報告
本文介紹了使用POI-TL和JFreeChart生成包含動態(tài)數(shù)據(jù)和圖表的Word報告的方法,并分享了實際開發(fā)中的踩坑經(jīng)驗,通過代碼示例講解的非常詳細(xì),具有一定的參考價值,需要的朋友可以參考下2025-02-02
探究springboot中的TomcatMetricsBinder
springboot的TomcatMetricsBinder主要是接收ApplicationStartedEvent然后創(chuàng)建TomcatMetrics執(zhí)行bindTo進行注冊,TomcatMetrics主要注冊了globalRequest、servlet、cache、threadPool、session相關(guān)的指標(biāo),本文給大家介紹的非常詳細(xì),需要的朋友參考下吧2023-11-11
Java easyExcel實現(xiàn)導(dǎo)入多sheet的Excel
這篇文章主要為大家詳細(xì)介紹了如何使用Java easyExcel實現(xiàn)導(dǎo)入多sheet的Excel,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以了解一下2025-06-06
Java實現(xiàn)游戲飛機大戰(zhàn)-III的示例代碼
這篇文章主要為大家介紹了如何利用Java實現(xiàn)經(jīng)典的游戲之飛機大戰(zhàn),文中采用了swing技術(shù)進行了界面化處理,感興趣的小伙伴可以動手試一試2022-02-02

