最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SpringCloud微服務(wù)開(kāi)發(fā)基于RocketMQ實(shí)現(xiàn)分布式事務(wù)管理詳解

 更新時(shí)間:2022年09月19日 14:54:48   作者:π大星的日常  
分布式事務(wù)是在微服務(wù)開(kāi)發(fā)中經(jīng)常會(huì)遇到的一個(gè)問(wèn)題,之前的文章中我們已經(jīng)實(shí)現(xiàn)了利用Seata來(lái)實(shí)現(xiàn)強(qiáng)一致性事務(wù),其實(shí)還有一種廣為人知的方案就是利用消息隊(duì)列來(lái)實(shí)現(xiàn)分布式事務(wù),保證數(shù)據(jù)的最終一致性,也就是我們常說(shuō)的柔性事務(wù)

消息隊(duì)列實(shí)現(xiàn)分布式事務(wù)原理

首先讓我們來(lái)看一下基于消息隊(duì)列實(shí)現(xiàn)分布式事務(wù)的原理方案。

柔性事務(wù)

發(fā)送消息的服務(wù)有個(gè)OUTBOX數(shù)據(jù)表,在進(jìn)行INSERT、UPDATE、DELETE 業(yè)務(wù)操作時(shí)也會(huì)給OUTBOX數(shù)據(jù)表INSERT一條消息記錄,這樣可以保證原子性,因?yàn)檫@是基于本地的ACID事務(wù)。

OUTBOX表充當(dāng)臨時(shí)消息隊(duì)列,然后我們?cè)谝胍粋€(gè)消息中繼(MessageRelay)的服務(wù),由他從OUTBOX表中讀取數(shù)據(jù)并發(fā)布消息到消息組件。

消息中繼的實(shí)現(xiàn)可以很簡(jiǎn)單,只需要通過(guò)定時(shí)任務(wù)定期從OUTBOX表中拉取最新未發(fā)布的數(shù)據(jù),獲取到數(shù)據(jù)后將數(shù)據(jù)發(fā)送給消息組件,最后將完成發(fā)送的消息從OUTBOX表中刪除即可,對(duì)于失敗的消息可以根據(jù)業(yè)務(wù)規(guī)則進(jìn)行重試。

RocketMQ的事務(wù)消息

RocketMQ本身已經(jīng)支持事務(wù)消息,如果你們項(xiàng)目使用了RocketMQ,可以直接借助RocketMQ的事務(wù)消息實(shí)現(xiàn)分布式事務(wù),我們先看一下RocketMQ事務(wù)消息的原理然后再借助RocketMQ來(lái)實(shí)現(xiàn)分布式事務(wù)。

RocketMQ采用了2PC的思想來(lái)實(shí)現(xiàn)了提交事務(wù)消息,同時(shí)增加一個(gè)補(bǔ)償邏輯來(lái)處理二階段超時(shí)或者失敗的消息,如下圖所示。

分布式事務(wù)

RocketMQ實(shí)現(xiàn)事務(wù)消息主要分為兩個(gè)階段:正常事務(wù)的發(fā)送及提交、事務(wù)信息的補(bǔ)償流程

整體流程為:

正常事務(wù)發(fā)送與提交階段

1、生產(chǎn)者發(fā)送一個(gè)半消息給MQServer(半消息是指消費(fèi)者暫時(shí)不能消費(fèi)的消息)

2、服務(wù)端響應(yīng)消息寫入結(jié)果,半消息發(fā)送成功

3、開(kāi)始執(zhí)行本地事務(wù)

4、根據(jù)本地事務(wù)的執(zhí)行狀態(tài)執(zhí)行Commit或者Rollback操作

事務(wù)信息的補(bǔ)償流程

1、如果MQServer長(zhǎng)時(shí)間沒(méi)收到本地事務(wù)的執(zhí)行狀態(tài)會(huì)向生產(chǎn)者發(fā)起一個(gè)確認(rèn)回查的操作請(qǐng)求

2、生產(chǎn)者收到確認(rèn)回查請(qǐng)求后,檢查本地事務(wù)的執(zhí)行狀態(tài)

3、根據(jù)檢查后的結(jié)果執(zhí)行Commit或者Rollback操作

補(bǔ)償階段主要是用于解決生產(chǎn)者在發(fā)送Commit或者Rollback操作時(shí)發(fā)生超時(shí)或失敗的情況。

RocketMQ事務(wù)流程關(guān)鍵

事務(wù)消息在一階段對(duì)用戶不可見(jiàn)

事務(wù)消息相對(duì)普通消息最大的特點(diǎn)就是一階段發(fā)送的消息對(duì)用戶是不可見(jiàn)的,也就是說(shuō)消費(fèi)者不能直接消費(fèi)。這里RocketMQ的實(shí)現(xiàn)方法是原消息的主題與消息消費(fèi)隊(duì)列,然后把主題改成RMQ_SYS_TRANS_HALF_TOPIC,這樣由于消費(fèi)者沒(méi)有訂閱這個(gè)主題,所以不會(huì)被消費(fèi)。

如何處理第二階段的失敗消息?

在本地事務(wù)執(zhí)行完成后會(huì)向MQServer發(fā)送Commit或Rollback操作,此時(shí)如果在發(fā)送消息的時(shí)候生產(chǎn)者出故障了,那么要保證這條消息最終被消費(fèi),MQServer會(huì)像服務(wù)端發(fā)送回查請(qǐng)求,確認(rèn)本地事務(wù)的執(zhí)行狀態(tài)。

當(dāng)然了rocketmq并不會(huì)無(wú)休止的的信息事務(wù)狀態(tài)回查,默認(rèn)回查15次,如果15次回查還是無(wú)法得知事務(wù)狀態(tài),RocketMQ默認(rèn)回滾該消息。

消息狀態(tài) 事務(wù)消息有三種狀態(tài):TransactionStatus.CommitTransaction:提交事務(wù)消息,消費(fèi)者可以消費(fèi)此消息

TransactionStatus.RollbackTransaction:回滾事務(wù),它代表該消息將被刪除,不允許被消費(fèi)。

TransactionStatus.Unknown:中間狀態(tài),它代表需要檢查消息隊(duì)列來(lái)確定狀態(tài)。

代碼實(shí)現(xiàn)

業(yè)務(wù)需求:用戶請(qǐng)求訂單微服務(wù)order-service接口刪除訂單(退貨),刪除訂單時(shí)需要調(diào)用account-service的方法給賬戶增加余額,一個(gè)典型的分布式事務(wù)問(wèn)題。

基礎(chǔ)配置

在Order-Service和Account-Service中引入Rocket消息組件

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
</dependency>

在配置中心添加RocketMQ的相關(guān)配置

rocketmq:
  name-server: xxx.xx.x.xx:9876
  producer:
    group: cloud-group

在OrderService服務(wù)中建立一張事務(wù)日志表rocketmq_transaction_log(作用稍后說(shuō))

發(fā)送半消息

Order-Service作為分布式事務(wù)開(kāi)始的入口,在Service層我們給RocketMQ發(fā)送一條半消息

OrderController入口

/**
 * 根據(jù)訂單號(hào)刪除訂單
 * @param orderNo 訂單編號(hào)
 */
@PostMapping("/order/delete")
public ResultData<String> delete(@RequestParam String orderNo){
 log.info("delete order id is {}",orderNo);
 orderService.delete(orderNo);
 return ResultData.success("訂單刪除成功");
}

直接調(diào)用orderService的delete方法

OrderServiceImpl業(yè)務(wù)邏輯

@Override
public void delete(String orderNo) {
 Order order = orderMapper.selectByNo(orderNo);
 //如果訂單存在且狀態(tài)為有效,進(jìn)行業(yè)務(wù)處理
 if (order != null && CloudConstant.VALID_STATUS.equals(order.getStatus())) {
  String transactionId = UUID.randomUUID().toString();
  //如果可以刪除訂單則發(fā)送消息給rocketmq,讓用戶中心消費(fèi)消息
  rocketMQTemplate.sendMessageInTransaction("add-amount",
    MessageBuilder.withPayload(
      UserAddMoneyDTO.builder()
        .userCode(order.getAccountCode())
        .amount(order.getAmount())
        .build()
    )
    .setHeader(RocketMQHeaders.TRANSACTION_ID, transactionId)
    .setHeader("order_id",order.getId())
    .build()
    ,order
  );
 
 }
}

首先校驗(yàn)一下訂單狀態(tài),然后使用rocketMQTemplate.sendMessageInTransaction()發(fā)送事務(wù)消息。

sendMessageInTransaction方法有三個(gè)參數(shù):

  • destination:目的地(主題),這里發(fā)送給add-amount這個(gè)topic
  • message:發(fā)送給消費(fèi)者的消息體,需要使用MessageBuilder.withPayload()來(lái)構(gòu)建消息
  • arg:參數(shù)

注意,這里我們生成了一個(gè)transactionId,并放在header中跟消息一起發(fā)送(這里實(shí)際也可以構(gòu)造成一個(gè)對(duì)象,放在arg里進(jìn)行發(fā)送),作用后面再講!

消息封裝實(shí)體UserAddMoneyDTO

@Data
@NoArgsConstructor
@AllArgsConstructor
@Builder
public class UserAddMoneyDTO {
    /**
     * 用戶編碼
     */
    private String userCode;
    /**
     * 金額
     */
    private BigDecimal amount;
}

這個(gè)類生產(chǎn)者和消費(fèi)者都需要用到,所以我直接丟到common包中,大家根據(jù)項(xiàng)目實(shí)際情況決定放哪。

執(zhí)行本地事務(wù)與回查

MQServer收到半消息后會(huì)告訴生產(chǎn)者order-service確認(rèn)收到半消息,這時(shí)候order-service需要執(zhí)行本地事務(wù),執(zhí)行完本地事務(wù)后再告訴MQServer本地事務(wù)的執(zhí)行狀態(tài),確認(rèn)此消息究竟是Commit還是Rollback。

RocketMQ提供了RocketMQLocalTransactionListener接口,本地事務(wù)監(jiān)聽(tīng)器,這個(gè)接口類的實(shí)現(xiàn)如下:

第一個(gè)方法executeLocalTransaction為執(zhí)行本地事務(wù);第二個(gè)方法checkLocalTransaction為檢查本地事務(wù)的執(zhí)行狀態(tài),也就是回查動(dòng)作。

我們需要實(shí)現(xiàn)RocketMQLocalTransactionListener接口,在executeLocalTransaction方法中執(zhí)行本地事務(wù),在執(zhí)行checkLocalTransaction回查方法時(shí)告訴RocketMQ到底該提交還是回滾。

這里大家思考一個(gè)問(wèn)題,本地事務(wù)已經(jīng)執(zhí)行完成了,怎么去回查本地事務(wù)的執(zhí)行結(jié)果呢?

答案如下:我們可以在執(zhí)行本地事務(wù)的時(shí)候同時(shí)生成一條事務(wù)日志,讓本地事務(wù)與日志事務(wù)在同一個(gè)方法中,同時(shí)添加@Transactional注解,保證兩個(gè)操作事務(wù)是一個(gè)原子操作。

這樣如果事務(wù)日志表中有這個(gè)本地事務(wù)的信息,那就代表本地事務(wù)執(zhí)行成功,需要Commit,相反如果沒(méi)有對(duì)應(yīng)的事務(wù)日志,則表示執(zhí)行失敗,需要Rollback。這就是為什么我們上面在OrderService中需要建立一張事務(wù)日志表的原因。

實(shí)現(xiàn)RocketMQLocalTransactionListener接口,完成事務(wù)執(zhí)行邏輯

/**
 * 監(jiān)聽(tīng)事務(wù)消息
 * @author javadaily
 */
@Slf4j
@RocketMQTransactionListener
@RequiredArgsConstructor(onConstructor = @__(@Autowired))
public class AddUserAmountListener implements RocketMQLocalTransactionListener {
    private final OrderService orderService;
    private final RocketMqTransactionLogMapper rocketMqTransactionLogMapper;
    /**
     * 執(zhí)行本地事務(wù)
     */
    @Override
    public RocketMQLocalTransactionState executeLocalTransaction(Message message, Object arg) {
        log.info("執(zhí)行本地事務(wù)");
        MessageHeaders headers = message.getHeaders();
        //獲取事務(wù)ID
        String transactionId = (String) headers.get(RocketMQHeaders.TRANSACTION_ID);
        Integer orderId = Integer.valueOf((String)headers.get("order_id"));
        log.info("transactionId is {}, orderId is {}",transactionId,orderId);
        try{
            //執(zhí)行本地事務(wù),并記錄日志
            orderService.changeStatuswithRocketMqLog(orderId, CloudConstant.INVALID_STATUS,transactionId);
            //執(zhí)行成功,可以提交事務(wù)
            return RocketMQLocalTransactionState.COMMIT;
        }catch (Exception e){
            return RocketMQLocalTransactionState.ROLLBACK;
        }
    }
    /**
     * 本地事務(wù)的檢查,檢查本地事務(wù)是否成功
     */
    @Override
    public RocketMQLocalTransactionState checkLocalTransaction(Message message) {
        MessageHeaders headers = message.getHeaders();
        //獲取事務(wù)ID
        String transactionId = (String) headers.get(RocketMQHeaders.TRANSACTION_ID);
        log.info("檢查本地事務(wù),事務(wù)ID:{}",transactionId);
        //根據(jù)事務(wù)id從日志表檢索
        QueryWrapper<RocketmqTransactionLog> queryWrapper = new QueryWrapper<>();
        queryWrapper.eq("transaction_id",transactionId);
        RocketmqTransactionLog rocketmqTransactionLog = rocketMqTransactionLogMapper.selectOne(queryWrapper);
        if(null != rocketmqTransactionLog){
            return RocketMQLocalTransactionState.COMMIT;
        }
        return RocketMQLocalTransactionState.ROLLBACK;
    }
}

本地事務(wù)執(zhí)行邏輯

@Transactional(rollbackFor = RuntimeException.class)
@Override
public void changeStatuswithRocketMqLog(Integer id,String status,String transactionId){
    orderMapper.changeStatus(id,status);
    rocketMqTransactionLogMapper.insert(
        RocketmqTransactionLog.builder()
        .transactionId(transactionId)
        .log("執(zhí)行刪除訂單操作")
        .build()
    );
}

修改訂單狀態(tài)為刪除狀態(tài),同時(shí)往事務(wù)日志表中插入一條事務(wù)日志,用@Transactional注解保證事務(wù)。

Account-Service消費(fèi)消息

監(jiān)聽(tīng)消息并處理給用戶增加余額邏輯

@Slf4j
@Service
@RocketMQMessageListener(topic = "add-amount",consumerGroup = "cloud-group")
@RequiredArgsConstructor(onConstructor = @__(@Autowired) )
public class AddUserAmountListener implements RocketMQListener<UserAddMoneyDTO> {
    private final AccountMapper accountMapper;
    /**
     * 收到消息的業(yè)務(wù)邏輯
     */
    @Override
    public void onMessage(UserAddMoneyDTO userAddMoneyDTO) {
        log.info("received message: {}",userAddMoneyDTO);
        accountMapper.increaseAmount(userAddMoneyDTO.getUserCode(),userAddMoneyDTO.getAmount());
        log.info("add money success");
    }
}

測(cè)試

測(cè)試數(shù)據(jù)

訂單表

用戶表

事務(wù)日志表

如果事務(wù)消息成功消費(fèi)最終用戶表中jianzh5這個(gè)用戶的amount應(yīng)該變成300(100+200)

測(cè)試準(zhǔn)備

我們?cè)趫?zhí)行本地事務(wù)成功并需要通知消息隊(duì)列提交事務(wù)處打個(gè)斷點(diǎn),然后在執(zhí)行到此處時(shí)手動(dòng)模擬異常

模擬異常

在準(zhǔn)備提交事務(wù)時(shí)我們通過(guò)命令taskkill /pid 10116 -t -f命令強(qiáng)制殺掉OrderService進(jìn)程。(先通過(guò)jps獲取OrderService進(jìn)程ID)

重啟服務(wù)器,檢查是否會(huì)執(zhí)行回查方法

重啟OrderService程序會(huì)自動(dòng)執(zhí)行回查方法,結(jié)合事務(wù)日志表判斷是否提交事務(wù)。

運(yùn)行后的結(jié)果

小結(jié)

我們介紹了使用消息隊(duì)列實(shí)現(xiàn)柔性事務(wù)的方案,重點(diǎn)剖析了RocketMQ事務(wù)消息的原理,并通過(guò)Demo案例實(shí)現(xiàn)了分布式事務(wù)(柔性事務(wù))。

到此這篇關(guān)于SpringCloud微服務(wù)開(kāi)發(fā)基于RocketMQ實(shí)現(xiàn)分布式事務(wù)管理詳解的文章就介紹到這了,更多相關(guān)SpringCloud RocketMQ分布式事務(wù)內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • SpringMVC HttpMessageConverter報(bào)文信息轉(zhuǎn)換器

    SpringMVC HttpMessageConverter報(bào)文信息轉(zhuǎn)換器

    ??HttpMessageConverter???,報(bào)文信息轉(zhuǎn)換器,將請(qǐng)求報(bào)文轉(zhuǎn)換為Java對(duì)象,或?qū)ava對(duì)象轉(zhuǎn)換為響應(yīng)報(bào)文。???HttpMessageConverter???提供了兩個(gè)注解和兩個(gè)類型:??@RequestBody,@ResponseBody???,??RequestEntity,ResponseEntity??
    2023-01-01
  • Java客戶端利用Jedis操作redis緩存示例代碼

    Java客戶端利用Jedis操作redis緩存示例代碼

    Jedis是Redis官方推薦的用于訪問(wèn)Java客戶端,下面這篇文章主要給大家介紹了關(guān)于Java客戶端利用Jedis操作redis緩存的相關(guān)資料,文中給出了詳細(xì)的示例代碼,需要的朋友可以參考借鑒,下面來(lái)一起看看吧。
    2017-07-07
  • 一文秒懂通過(guò)JavaCSV類庫(kù)讀寫CSV文件的技巧

    一文秒懂通過(guò)JavaCSV類庫(kù)讀寫CSV文件的技巧

    本文給大家推薦第三方工具庫(kù) JavaCSV,用來(lái)造一些 csv 測(cè)試數(shù)據(jù)文件,使用超級(jí)方便,本文通過(guò)示例代碼給大家介紹的非常詳細(xì),感興趣的朋友跟隨小編一起看看吧
    2021-05-05
  • Java concurrency之共享鎖和ReentrantReadWriteLock_動(dòng)力節(jié)點(diǎn)Java學(xué)院整理

    Java concurrency之共享鎖和ReentrantReadWriteLock_動(dòng)力節(jié)點(diǎn)Java學(xué)院整理

    本篇文章主要介紹了Java concurrency之共享鎖和ReentrantReadWriteLock,非常具有實(shí)用價(jià)值,需要的朋友可以參考下
    2017-06-06
  • SpringBoot查詢數(shù)據(jù)庫(kù)導(dǎo)出報(bào)表文件方式

    SpringBoot查詢數(shù)據(jù)庫(kù)導(dǎo)出報(bào)表文件方式

    這篇文章主要介紹了SpringBoot查詢數(shù)據(jù)庫(kù)導(dǎo)出報(bào)表文件方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2024-04-04
  • Spring中最常用的注解之一@Autowired詳解

    Spring中最常用的注解之一@Autowired詳解

    本文講解了Spring中最常用的注解之一@Autowired, 平時(shí)我們可能都是使用屬性注入的,但是后續(xù)建議大家慢慢改變習(xí)慣,使用構(gòu)造器注入。同時(shí)也講解了這個(gè)注解背后的實(shí)現(xiàn)原理,需要的朋友可以參考下
    2023-01-01
  • 深入淺出Java中重試機(jī)制的多種方式

    深入淺出Java中重試機(jī)制的多種方式

    重試機(jī)制在分布式系統(tǒng)中,或者調(diào)用外部接口中,都是十分重要的。重試機(jī)制可以保護(hù)系統(tǒng)減少因網(wǎng)絡(luò)波動(dòng)、依賴服務(wù)短暫性不可用帶來(lái)的影響,讓系統(tǒng)能更穩(wěn)定的運(yùn)行的一種保護(hù)機(jī)制。本文就來(lái)和大家聊聊Java中重試機(jī)制的多種方式
    2023-03-03
  • 解決idea npm:無(wú)法將“npm”項(xiàng)識(shí)別為cmdlet、函數(shù)、腳本文件或可運(yùn)行程序的名稱問(wèn)題

    解決idea npm:無(wú)法將“npm”項(xiàng)識(shí)別為cmdlet、函數(shù)、腳本文件或可運(yùn)行程序的名稱問(wèn)題

    在IDEA中運(yùn)行npm命令時(shí)出現(xiàn)無(wú)法識(shí)別的錯(cuò)誤,通常是由于npm環(huán)境變量配置不正確引起,解決方法包括以管理員身份運(yùn)行IDEA,確認(rèn)node和npm是否正確安裝及配置環(huán)境變量,需要在系統(tǒng)環(huán)境變量中添加node.js的安裝路徑,并設(shè)置npm的全局模塊和緩存路徑
    2024-10-10
  • 淺談Java由于不當(dāng)?shù)膱?zhí)行順序?qū)е碌乃梨i

    淺談Java由于不當(dāng)?shù)膱?zhí)行順序?qū)е碌乃梨i

    為了保證線程的安全,我們引入了加鎖機(jī)制,但是如果不加限制的使用加鎖,就有可能會(huì)導(dǎo)致順序死鎖(Lock-Ordering Deadlock)。本文將會(huì)討論一下順序死鎖的問(wèn)題。
    2021-06-06
  • Spring Boot 框架詳細(xì)指南

    Spring Boot 框架詳細(xì)指南

    Spring Boot 是由 Pivotal 團(tuán)隊(duì)開(kāi)發(fā)的一個(gè)開(kāi)源 Java 框架,旨在簡(jiǎn)化 Spring 應(yīng)用程序的創(chuàng)建和開(kāi)發(fā)過(guò)程,這篇文章主要介紹了Spring Boot 框架詳細(xì)指南,需要的朋友可以參考下
    2025-05-05

最新評(píng)論

泗阳县| 河南省| 武乡县| 榆中县| 普安县| 大兴区| 开远市| 温州市| 太仓市| 南和县| 阆中市| 龙川县| 宁晋县| 桦川县| 德庆县| 北碚区| 内丘县| 内丘县| 肃宁县| 丹凤县| 长岛县| 罗山县| 镇赉县| 昌江| 利川市| 浦东新区| 侯马市| 荥阳市| 鄄城县| 黑山县| 长兴县| 凌源市| 两当县| 安达市| 静宁县| 甘南县| 泰顺县| 龙门县| 嘉定区| 凉山| 凤阳县|