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

RocketMQ事務(wù)消息機(jī)制詳解

 更新時(shí)間:2024年01月11日 09:28:50   作者:智由靜生  
這篇文章主要介紹了RocketMQ事務(wù)消息機(jī)制詳解,RocketMQ服務(wù)端將消息持久化之后,向發(fā)送方返回Ack確認(rèn)消息已經(jīng)發(fā)送成功,由于消息為半事務(wù)消息,在未收到生產(chǎn)者對(duì)該消息的二次確認(rèn)前,此消息被標(biāo)記成"暫不能投遞"狀態(tài),需要的朋友可以參考下

RocketMQ事務(wù)消息

RocketMQ提供了事務(wù)消息,通過(guò)事務(wù)消息就能達(dá)到分布式事務(wù)的最終一致,從而實(shí)現(xiàn)了可靠消息服務(wù)。

一、事務(wù)消息的實(shí)現(xiàn)步驟

事務(wù)消息發(fā)送步驟:

1. 發(fā)送方將半事務(wù)消息發(fā)送至RocketMQ服務(wù)端。

2. RocketMQ服務(wù)端將消息持久化之后,向發(fā)送方返回Ack確認(rèn)消息已經(jīng)發(fā)送成功。由于消息為半事務(wù)消息,在未收到生產(chǎn)者對(duì)該消息的二次確認(rèn)前,此消息被標(biāo)記成“暫不能投遞”狀態(tài)。

3. 發(fā)送方開(kāi)始執(zhí)行本地事務(wù)邏輯。

4. 發(fā)送方根據(jù)本地事務(wù)執(zhí)行結(jié)果向服務(wù)端提交二次確認(rèn)(Commit 或是 Rollback),服務(wù)端收到Commit 狀態(tài)則將半事務(wù)消息標(biāo)記為可投遞,訂閱方最終將收到該消息;服務(wù)端收到 Rollback 狀態(tài)則刪除半事務(wù)消息,訂閱方將不會(huì)接受該消息。

事務(wù)消息回查步驟:

1. 在斷網(wǎng)或者是應(yīng)用重啟的特殊情況下,上述步驟4提交的二次確認(rèn)最終未到達(dá)服務(wù)端,經(jīng)過(guò)固定時(shí)間后服務(wù)端將對(duì)該消息發(fā)起消息回查。

2. 發(fā)送方收到消息回查后,需要檢查對(duì)應(yīng)消息的本地事務(wù)執(zhí)行的最終結(jié)果。 3. 發(fā)送方根據(jù)檢查得到的本地事務(wù)的最終狀態(tài)再次提交二次確認(rèn),服務(wù)端仍按照步驟4對(duì)半事務(wù)消息進(jìn)行操作。

二、程序?qū)崿F(xiàn)

事務(wù)消息處理類需要繼承RocketMQLocalTransactionListener類。該類的executeLocalTransaction方法負(fù)責(zé)在接到RocketMQ服務(wù)端的Ack確認(rèn)消息后執(zhí)行本地方法,也就是事務(wù)消息發(fā)送步驟中的步驟3。該類的checkLocalTransaction方法負(fù)責(zé),在斷網(wǎng)或者是應(yīng)用重啟的特殊情況下,執(zhí)行RocketMQ服務(wù)端的消息回查,也就是事務(wù)消息回查步驟中的步驟2。

此外,要使該類生效,還需要加@RocketMQTransactionListener注解。這里有個(gè)要特別注意的地方。在2.1.0版本前,這個(gè)注解有一個(gè)屬性txProducerGroup,可以用多個(gè)@RocketMQTransactionListener來(lái)監(jiān)聽(tīng)不同的txProducerGroup來(lái)發(fā)送不同類型的事務(wù)消息到topic。但是現(xiàn)在在一個(gè)項(xiàng)目中,如果你在一個(gè)project中寫(xiě)了多個(gè)@RocketMQTransactionListener,項(xiàng)目將不能啟動(dòng),啟動(dòng)會(huì)報(bào)錯(cuò)。產(chǎn)生這個(gè)問(wèn)題的原因據(jù)說(shuō)是,當(dāng)使用RocketMQTemplate并發(fā)的執(zhí)行事務(wù)時(shí),非常容易出現(xiàn)"illegal state"的異常,原因是一個(gè)TransactionProducer在執(zhí)行事務(wù)時(shí)不能被共享。所以,必須使用同一個(gè)TransactionMQProducer來(lái)發(fā)送所有類型的事務(wù)消息。當(dāng)然同理也就必須使用一個(gè)偵聽(tīng)器處理所有的消息了。

既然必須使用同一個(gè)TransactionMQProducer,對(duì)于比較大的應(yīng)用,業(yè)務(wù)場(chǎng)景很多,就會(huì)造成混亂。這里我給出一個(gè)方案拋磚引玉。TransactionMQProducer在發(fā)送消息時(shí),是可以傳遞參數(shù)對(duì)象和指定消息頭的??梢园岩獔?zhí)行的本地方法的bean名和方法名放進(jìn)去。

//發(fā)送半事務(wù)消息
TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
		topicAndTag,
		MessageBuilder.withPayload(msg)
			.setHeader(Constants.TX_ID_HEADER_NAME, msg.getTxId())
			.setHeader(Constants.CHECK_BEAN_ID_HEADER_NAME, def.getCheckBeanId())
			.setHeader(Constants.BIZ_ID_HEADER_NAME, msg.getBizId())
			.build(),
		def
);

其中def就是參數(shù)對(duì)象,可以自定義對(duì)象,這里是我自定義的TransactionMsgDefinationDto類,可以把想傳遞的信息放進(jìn)去,最重要的是要執(zhí)行的本地方法的bean名和方法名和方法執(zhí)行參數(shù):executeBeanId(bean名)、executeBeanMethod(方法名)、executeBeanParams(方法執(zhí)行參數(shù))。該對(duì)象可以傳給RocketMQLocalTransactionListener的executeLocalTransaction方法,然后通過(guò)反射執(zhí)行。

@Override
	public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
		try {
			//保存消息記錄
			String body = new String((byte[]) msg.getPayload(), StandardCharsets.UTF_8);
			JSONObject jsonBody = JSONObject.parseObject(body);
			BaseMsgDto dto = JSONObject.toJavaObject(jsonBody, BaseMsgDto.class);//(BaseMsgDto)msg.getPayload();
			TransactionMsgDefinationDto def = (TransactionMsgDefinationDto)arg;
			ProducerLog producerLog = BeanCopyUtils.copyProperties(def, ProducerLog::new);
			String[] tags = def.getMsgTags();
			if(tags !=null && tags.length > 0) {
				StringBuilder tag = new StringBuilder();
				for(int i = 0; i<tags.length; i++) {
					tag.append(tags[0]);
					if(i != tags.length-1) {
						tag.append("||");
					}
				}
				producerLog.setMsgTag(tag.toString());
			}
			producerLog.setBizId(dto.getBizId());
			producerLog.setTxId(dto.getTxId());
			producerLog.setBizType(dto.getBizType());
			producerLog.setGroupName(dto.getProducerGroup());
			producerLog.setMsgBody(body);
			producerLogService.save(producerLog);
			//執(zhí)行事務(wù)方法
			SpringUtil.invokeBeanMethod(def.getExecuteBeanId(), def.getExecuteBeanMethod(), def.getExecuteBeanParams());
			return RocketMQLocalTransactionState.COMMIT;
		} catch (Exception e) {
			logger.error("發(fā)生錯(cuò)誤:", e);
			return RocketMQLocalTransactionState.UNKNOWN;
		}
	}

放在消息頭header中的數(shù)據(jù)可以傳遞給RocketMQLocalTransactionListener的checkLocalTransaction方法,然后同樣通過(guò)反射執(zhí)行。

@Override
	public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
		try {
			String txId = (String)msg.getHeaders().get(Constants.TX_ID_HEADER_NAME); 
			String checkBeanId = (String)msg.getHeaders().get(Constants.CHECK_BEAN_ID_HEADER_NAME);
			Long bizId = Long.parseLong((String)msg.getHeaders().get(Constants.BIZ_ID_HEADER_NAME));
			//執(zhí)行檢查方法
			Boolean ret = (Boolean)SpringUtil.invokeBeanMethod(checkBeanId, "check", new Object[]{bizId, txId});
			if(ret.booleanValue())
				return RocketMQLocalTransactionState.COMMIT;
			else
				return RocketMQLocalTransactionState.ROLLBACK;
		} catch (Exception e) {
			logger.error("發(fā)生錯(cuò)誤:", e);
			return RocketMQLocalTransactionState.UNKNOWN;
		}
	}

到此這篇關(guān)于RocketMQ事務(wù)消息機(jī)制詳解的文章就介紹到這了,更多相關(guān)RocketMQ事務(wù)消息內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • springboot 項(xiàng)目使用jasypt加密數(shù)據(jù)源的方法

    springboot 項(xiàng)目使用jasypt加密數(shù)據(jù)源的方法

    Jasypt 是一個(gè) Java 庫(kù),它允許開(kāi)發(fā)者以最小的努力為他/她的項(xiàng)目添加基本的加密功能,而且不需要對(duì)密碼學(xué)的工作原理有深刻的了解。接下來(lái)通過(guò)本文給大家介紹springboot 項(xiàng)目使用jasypt加密數(shù)據(jù)源的問(wèn)題,一起看看吧
    2021-11-11
  • Spring Boot使用過(guò)濾器和攔截器分別實(shí)現(xiàn)REST接口簡(jiǎn)易安全認(rèn)證示例代碼詳解

    Spring Boot使用過(guò)濾器和攔截器分別實(shí)現(xiàn)REST接口簡(jiǎn)易安全認(rèn)證示例代碼詳解

    這篇文章主要介紹了Spring Boot使用過(guò)濾器和攔截器分別實(shí)現(xiàn)REST接口簡(jiǎn)易安全認(rèn)證示例代碼,通過(guò)開(kāi)發(fā)實(shí)踐,理解過(guò)濾器和攔截器的工作原理,需要的朋友可以參考下
    2018-06-06
  • SpringBoot3整合SpringSecurity6快速入門示例教程

    SpringBoot3整合SpringSecurity6快速入門示例教程

    SpringSecurity 是Spring大家族中一名重要成員,是專門負(fù)責(zé)安全的框架,本文給大家介紹SpringBoot3整合SpringSecurity6快速入門示例教程,感興趣的朋友一起看看吧
    2025-04-04
  • Java調(diào)用打印機(jī)的2種方式舉例(無(wú)驅(qū)/有驅(qū))

    Java調(diào)用打印機(jī)的2種方式舉例(無(wú)驅(qū)/有驅(qū))

    我們平時(shí)使用某些軟件或者在超市購(gòu)物的時(shí)候都會(huì)發(fā)現(xiàn)可以使用打印機(jī)進(jìn)行打印,這篇文章主要給大家介紹了關(guān)于Java調(diào)用打印機(jī)的2種方式,分別是無(wú)驅(qū)/有驅(qū)的相關(guān)資料,需要的朋友可以參考下
    2023-11-11
  • idea創(chuàng)建xml文件全過(guò)程

    idea創(chuàng)建xml文件全過(guò)程

    總結(jié):通過(guò)File->Settings->Editor->FileAndCodeTemplates,創(chuàng)建一個(gè)自定義的XML文件模板,并命名為XMLFile.xml,后綴名為xml,模板內(nèi)容可自定義,并啟用實(shí)時(shí)模板功能,然后在文件夾中右鍵New,即可找到并創(chuàng)建XML文件
    2025-11-11
  • 使用Java把文本內(nèi)容轉(zhuǎn)換成網(wǎng)頁(yè)的實(shí)現(xiàn)方法分享

    使用Java把文本內(nèi)容轉(zhuǎn)換成網(wǎng)頁(yè)的實(shí)現(xiàn)方法分享

    這篇文章主要介紹了使用Java把文本內(nèi)容轉(zhuǎn)換成網(wǎng)頁(yè)的實(shí)現(xiàn)方法分享,利用到了Java中的文件io包,需要的朋友可以參考下
    2015-11-11
  • mybatis-plus雪花算法生成Id使用詳解

    mybatis-plus雪花算法生成Id使用詳解

    本文主要介紹了mybatis-plus雪花算法生成Id使用詳解,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2022-07-07
  • 詳解Spring AOP 實(shí)現(xiàn)主從讀寫(xiě)分離

    詳解Spring AOP 實(shí)現(xiàn)主從讀寫(xiě)分離

    本篇文章主要介紹了Spring AOP 實(shí)現(xiàn)主從讀寫(xiě)分離,小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2017-03-03
  • java導(dǎo)出dbf文件生僻漢字處理方式

    java導(dǎo)出dbf文件生僻漢字處理方式

    這篇文章主要介紹了java導(dǎo)出dbf文件生僻漢字處理方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-06-06
  • 詳解springboot WebTestClient的使用

    詳解springboot WebTestClient的使用

    WebClient是一個(gè)響應(yīng)式客戶端,它提供了RestTemplate的替代方法。這篇文章主要介紹了詳解springboot WebTestClient的使用, 具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2018-11-11

最新評(píng)論

海伦市| 沂水县| 铜梁县| 闻喜县| 神池县| 时尚| 桃源县| 清远市| 武功县| 本溪市| 海晏县| 台中市| 佛冈县| 太仓市| 满洲里市| 广河县| 庄浪县| 石台县| 江华| 团风县| 西乌珠穆沁旗| 洛川县| 陆丰市| 赣州市| 碌曲县| 永济市| 松溪县| 南溪县| 安达市| 安化县| 渭源县| 涞水县| 九江县| 昂仁县| 太保市| 永和县| 万年县| 茂名市| 房山区| 兴化市| 孝昌县|