一区二区三区在线-一区二区三区亚洲视频-一区二区三区亚洲-一区二区三区午夜-一区二区三区四区在线视频-一区二区三区四区在线免费观看

服務器之家:專注于服務器技術及軟件下載分享
分類導航

PHP教程|ASP.NET教程|Java教程|ASP教程|編程技術|正則表達式|C/C++|IOS|C#|Swift|Android|VB|R語言|JavaScript|易語言|vb.net|

服務器之家 - 編程語言 - Java教程 - 微服務架構設計RocketMQ進階事務消息原理詳解

微服務架構設計RocketMQ進階事務消息原理詳解

2022-03-04 00:37飄渺Jam Java教程

這篇文章主要介紹了為大家介紹了微服務架構中RocketMQ進階層面事務消息的原理詳解,有需要的朋友可以借鑒參考下希望能夠有所幫助

前言

分布式消息選型的時候是否支持事務消息是一個很重要的考量點,而目前只有RocketMQ對事務消息支持的最好。今天我們來嘮嘮如何實現RocketMQ的事務消息!

Apache RocketMQ在4.3.0版中已經支持分布式事務消息,這里RocketMQ采用了2PC的思想來實現了提交事務消息,同時增加一個補償邏輯來處理二階段超時或者失敗的消息,如下圖所示。

微服務架構設計RocketMQ進階事務消息原理詳解

 

RocketMQ事務流程概要

RocketMQ實現事務消息主要分為兩個階段:正常事務的發送及提交、事務信息的補償流程 整體流程為:

正常事務發送與提交階段
1、生產者發送一個半消息給MQServer(半消息是指消費者暫時不能消費的消息)
2、服務端響應消息寫入結果,半消息發送成功
3、開始執行本地事務
4、根據本地事務的執行狀態執行Commit或者Rollback操作

事務信息的補償流程
1、如果MQServer長時間沒收到本地事務的執行狀態會向生產者發起一個確認回查的操作請求
2、生產者收到確認回查請求后,檢查本地事務的執行狀態
3、根據檢查后的結果執行Commit或者Rollback操作
補償階段主要是用于解決生產者在發送Commit或者Rollback操作時發生超時或失敗的情況。

 

RocketMQ事務流程關鍵

1、事務消息在一階段對用戶不可見
事務消息相對普通消息最大的特點就是一階段發送的消息對用戶是不可見的,也就是說消費者不能直接消費。這里RocketMQ的實現方法是原消息的主題與消息消費隊列,然后把主題改成RMQ_SYS_TRANS_HALF_TOPIC ,這樣由于消費者沒有訂閱這個主題,所以不會被消費。

2、如何處理第二階段的失敗消息?
在本地事務執行完成后會向MQServer發送Commit或Rollback操作,此時如果在發送消息的時候生產者出故障了,那么要保證這條消息最終被消費,MQServer會像服務端發送回查請求,確認本地事務的執行狀態。
當然了rocketmq并不會無休止的的信息事務狀態回查,默認回查15次,如果15次回查還是無法得知事務狀態,RocketMQ默認回滾該消息。

3、消息狀態 事務消息有三種狀態:

TransactionStatus.CommitTransaction:提交事務消息,消費者可以消費此消息

TransactionStatus.RollbackTransaction:回滾事務,它代表該消息將被刪除,不允許被消費。

TransactionStatus.Unknown :中間狀態,它代表需要檢查消息隊列來確定狀態。

 

實現

我們構建這樣一個需求:用戶請求訂單微服務 order-service 接口刪除訂單(退貨),刪除訂單后需要發送消息給用戶服務account-service,用戶微服務收到消息后會給用戶賬戶增加余額。這個需求跟錢相關,肯定要保證消息的事務性,接下來我們根據上面的原理實現整個流程。

基礎配置

生產者order-servcie和account-service都要引入RocketMQ相關依賴,增加RocketMQ的相關配置

引入組件

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

添加配置

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

發送半消息

order-service在執行刪除訂單操作時發送一條半消息給MQServer,發送半消息主要是使用rocketMQTemplate.sendMessageInTransaction() 方法,發送事務消息。

@Override
public void delete(String orderNo) {
	Order order = orderMapper.selectByNo(orderNo);
	//如果訂單存在且狀態為有效,進行業務處理
	if (order != null && CloudConstant.VALID_STATUS.equals(order.getStatus())) {
		String transactionId = UUID.randomUUID().toString();
		//如果可以刪除訂單則發送消息給rocketmq,讓用戶中心消費消息
		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
		);
	}
}

首先先校驗一下訂單狀態,然后發送消息給MQServer,這個邏輯大家都看得懂,

主要是關注sendMessageInTransaction() 方法,源碼如下:

public TransactionSendResult sendMessageInTransaction(String destination, Message<?> message, Object arg) throws MessagingException {
	try {
		if (((TransactionMQProducer)this.producer).getTransactionListener() == null) {
			throw new IllegalStateException("The rocketMQTemplate does not exist TransactionListener");
		} else {
			org.apache.rocketmq.common.message.Message rocketMsg = this.createRocketMqMessage(destination, message);
			return this.producer.sendMessageInTransaction(rocketMsg, arg);
		}
	} catch (MQClientException var5) {
		throw RocketMQUtil.convert(var5);
	}
}

該方法有三個參數:

destination:目的地(主題),這里發送給add-amount 這個主題

message:發送給消費者的消息體,需要使用 MessageBuilder.withPayload() 來構建消息

arg:參數

注意,這里我們生成了一個transactionId,并放在header中跟消息一起發送(這里實際也可以構造成一個對象,放在arg里進行發送),作用后面再講!

執行本地事務與回查

MQServer收到半消息后會告訴生產者order-service確認收到半消息,這時候order-service需要執行本地事務,執行完本地事務后再告訴MQServer本地事務的執行狀態,確認消息究竟是Commit還是Rollback。如果在告訴MQServer本地執行狀態的時候出異常了還需要讓MQServer能夠回查到,怎么實現這一些列操作呢?

RocketMQ提供了 RocketMQLocalTransactionListener 接口,本地事務監聽器,這個接口類的實現如下:

微服務架構設計RocketMQ進階事務消息原理詳解

第一個方法executeLocalTransaction 為執行本地事務;
第二個方法checkLocalTransaction 為檢查本地事務的執行狀態,也就是回查動作。
有了這個接口類我們的執行邏輯清楚了,但是還有個問題:本地事務已經執行完成了,怎么去回查本地事務的執行結果呢?

我們可以在執行本地事務的時候同時生成一個事務日志,讓本地事務與日志事務在同一個方法中,同時添加@Transactional 注解,保證兩個操作事務是一個原子操作。這樣如果事務日志表中有這個本地事務的信息,那就代表本地事務執行成功,需要Commit,相反如果沒有對應的事務日志,則表示沒執行成功,需要Rollback

思路既然理順了,咱們就開擼。

首先創建一個日志表

微服務架構設計RocketMQ進階事務消息原理詳解

很簡單的三個字段,主要是這個事務id,需要根據這個事務id回查事務,還記得我們在發送半消息時生成的事務id嗎,就是干這個用的!

在生產者編寫方法實現RocketMQLocalTransactionListener

@Slf4j
@RocketMQTransactionListener
@RequiredArgsConstructor(onConstructor = @__(@Autowired))
public class AddUserAmountListener implements RocketMQLocalTransactionListener {
  private final OrderService orderService;
  private final RocketMqTransactionLogMapper rocketMqTransactionLogMapper;
  /**
   * 執行本地事務
   */
  @Override
  public RocketMQLocalTransactionState executeLocalTransaction(Message message, Object arg) {
      log.info("執行本地事務");
      MessageHeaders headers = message.getHeaders();
      //獲取事務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{
          //執行本地事務,并記錄日志
          orderService.changeStatuswithRocketMqLog(orderId, CloudConstant.INVALID_STATUS,transactionId);
          //執行成功,可以提交事務
          return RocketMQLocalTransactionState.COMMIT;
      }catch (Exception e){
          return RocketMQLocalTransactionState.ROLLBACK;
      }
  }

  /**
   * 本地事務的檢查,檢查本地事務是否成功
   */
  @Override
  public RocketMQLocalTransactionState checkLocalTransaction(Message message) {

      MessageHeaders headers = message.getHeaders();
      //獲取事務ID
      String transactionId = (String) headers.get(RocketMQHeaders.TRANSACTION_ID);
      log.info("檢查本地事務,事務ID:{}",transactionId);
      //根據事務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;
  }
}

執行本地事務的方法

@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("執行刪除訂單操作")
			.build()
	);
}

這一塊的代碼邏輯都是在生產端,即Order-Server,大家不要搞錯了

消費消息

Rollback的消息MQServer會給我們處理,我們只要關注Commit狀態時消費端可以正常消費即可。在account-service監聽消息,如果收到消息則給用戶賬戶增加余額。

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

 

測試

微服務架構設計RocketMQ進階事務消息原理詳解

訂單表有這樣一條記錄,用戶為jianzh5,amount為200

微服務架構設計RocketMQ進階事務消息原理詳解

用戶表的記錄,執行完成后jianzh5的賬戶應該變成250

調用刪除訂單接口,刪除訂單

微服務架構設計RocketMQ進階事務消息原理詳解

發送半消息

微服務架構設計RocketMQ進階事務消息原理詳解

執行本地事務,并生成事務日志

微服務架構設計RocketMQ進階事務消息原理詳解

模擬異常情況 在發送Commit消息的時候我們用命令殺掉進程taskkill /pid 19748 -t -f,模擬異常!

微服務架構設計RocketMQ進階事務消息原理詳解

重新啟動order-service,查看是否會執行回查動作

微服務架構設計RocketMQ進階事務消息原理詳解

MQServer進行回查,檢查事務日志,判斷是否可以提交事務

消費者消費事務消息,保證事務的一致性

微服務架構設計RocketMQ進階事務消息原理詳解

 

總結

使用RocketMQ實現事務消息的過程還是很復雜的,需要好好理解開頭的那張圖,只有理解了事務消息的交互過程才能編寫相應的代碼!

好了,各位朋友們,本期的內容到此就全部結束啦,能看到這里的同學都是優秀的同學,下一個升職加薪的就是你了!

以上就是微服務架構RocketMQ進階事務消息原理詳解的詳細內容,更多關于微服務架構RocketMQ事務消息的資料請關注服務器之家其它相關文章!

原文鏈接:https://jianzh5.blog.csdn.net/article/details/105283647

延伸 · 閱讀

精彩推薦
主站蜘蛛池模板: 亚洲国产精品第一区二区三区 | 国产中文字幕 | 国产一区二区视频在线 | 我的男友是消防员在线观看 | 奇米狠狠色 | 国产亚洲精品九九久在线观看 | 福利一区二区在线观看 | 久久综合久综合久久鬼色 | 国产一区二区免费在线 | 欧美xingai| 国产精品一区二区三区久久 | 欧美午夜视频一区二区 | 韩国三级在线高速影院 | 国产成人免费高清激情视频 | 色综合合久久天天综合绕视看 | 图片专区小说专区卡通动漫 | 亚洲免费小视频 | 五月天国产精品 | 午夜无码国产理论在线 | 美女班主任让我爽了一夜视频 | 免费稚嫩福利 | 完整秽淫刺激长篇小说 | 天堂网在线网站成人午夜网站 | 小SAO货叫大声点妓女 | 天堂网在线.www天堂在线视频 | 国产精品久久久久久搜索 | sss亚洲国产欧美一区二区 | avtt手机版 | 俄罗斯freeⅹ性欧美 | 色老板视频在线观看 | 9 1 视频在线 | 欧美黑人换爱交换乱理伦片 | 消息称老熟妇乱视频一区二区 | 亚洲swag精品自拍一区 | 亚洲网站在线看 | 久久性综合亚洲精品电影网 | 国产午夜亚洲精品理论片不卡 | 高清一级片 | 国产亚洲玖玖玖在线观看 | 男人女人日皮视频 | 牛牛色婷婷在线视频播放 |