日韩性视频-久久久蜜桃-www中文字幕-在线中文字幕av-亚洲欧美一区二区三区四区-撸久久-香蕉视频一区-久久无码精品丰满人妻-国产高潮av-激情福利社-日韩av网址大全-国产精品久久999-日本五十路在线-性欧美在线-久久99精品波多结衣一区-男女午夜免费视频-黑人极品ⅴideos精品欧美棵-人人妻人人澡人人爽精品欧美一区-日韩一区在线看-欧美a级在线免费观看

歡迎訪問 生活随笔!

生活随笔

當前位置: 首頁 > 编程资源 > 编程问答 >内容正文

编程问答

事务消息的实现

發布時間:2024/4/13 编程问答 17 豆豆
生活随笔 收集整理的這篇文章主要介紹了 事务消息的实现 小編覺得挺不錯的,現在分享給大家,幫大家做個參考.

發送事務消息

1) 創建事務性生產者

使用 TransactionMQProducer類創建生產者,并指定唯一的 ProducerGroup,就可以設置自定義線程池來處理這些檢查請求。執行本地事務后、需要根據執行結果對消息隊列進行回復。回傳的事務狀態在請參考前一節。

package com.leon.mq.rocketmq.transaction;import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.client.producer.*; import org.apache.rocketmq.common.message.Message; import org.apache.rocketmq.common.message.MessageExt;import java.util.concurrent.TimeUnit;/*** 發送同步消息*/ public class Producer {public static void main(String[] args) throws Exception {//1.創建消息生產者producer,并制定生產者組名TransactionMQProducer producer = new TransactionMQProducer("group5");//2.指定Nameserver地址producer.setNamesrvAddr("192.168.25.135:9876;192.168.25.138:9876");//添加事務監聽器producer.setTransactionListener(new TransactionListener() {/*** 在該方法中執行本地事務* @param msg* @param arg* @return*/@Overridepublic LocalTransactionState executeLocalTransaction(Message msg, Object arg) {if (StringUtils.equals("TAGA", msg.getTags())) {return LocalTransactionState.COMMIT_MESSAGE;} else if (StringUtils.equals("TAGB", msg.getTags())) {return LocalTransactionState.ROLLBACK_MESSAGE;} else if (StringUtils.equals("TAGC", msg.getTags())) {return LocalTransactionState.UNKNOW;}return LocalTransactionState.UNKNOW;}/*** 該方法時MQ進行消息事務狀態回查* @param msg* @return*/@Overridepublic LocalTransactionState checkLocalTransaction(MessageExt msg) {System.out.println("消息的Tag:" + msg.getTags());return LocalTransactionState.COMMIT_MESSAGE;}});//3.啟動producerproducer.start();String[] tags = {"TAGA", "TAGB", "TAGC"};for (int i = 0; i < 3; i++) {//4.創建消息對象,指定主題Topic、Tag和消息體/*** 參數一:消息主題Topic* 參數二:消息Tag* 參數三:消息內容*/Message msg = new Message("TransactionTopic", tags[i], ("Hello World" + i).getBytes());//5.發送消息SendResult result = producer.sendMessageInTransaction(msg, null);//發送狀態SendStatus status = result.getSendStatus();System.out.println("發送結果:" + result);//線程睡1秒TimeUnit.SECONDS.sleep(2);}//6.關閉生產者producer//producer.shutdown();} }

2)實現事務的監聽接口

當發送半消息成功時,我們使用 executeLocalTransaction 方法來執行本地事務。它返回前一節中提到的三個事務狀態之一。checkLocalTranscation 方法用于檢查本地事務狀態,并回應消息隊列的檢查請求。它也是返回前一節中提到的三個事務狀態之一。

package com.leon.mq.rocketmq.transaction;import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently; import org.apache.rocketmq.common.message.MessageExt;import java.util.List;/*** 消息的接受者*/ public class Consumer {public static void main(String[] args) throws Exception {//1.創建消費者Consumer,制定消費者組名DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("group1");//2.指定Nameserver地址consumer.setNamesrvAddr("192.168.25.135:9876;192.168.25.138:9876");//3.訂閱主題Topic和Tagconsumer.subscribe("TransactionTopic", "*");//4.設置回調函數,處理消息consumer.registerMessageListener(new MessageListenerConcurrently() {//接受消息內容@Overridepublic ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {for (MessageExt msg : msgs) {System.out.println("consumeThread=" + Thread.currentThread().getName() + "," + new String(msg.getBody()));}return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;}});//5.啟動消費者consumerconsumer.start();System.out.println("生產者啟動");} }

使用限制

  • 事務消息不支持延時消息和批量消息。

  • 為了避免單個消息被檢查太多次而導致半隊列消息累積,我們默認將單個消息的檢查次數限制為 15 次,但是用戶可以通過 Broker 配置文件的 transactionCheckMax參數來修改此限制。如果已經檢查某條消息超過 N 次的話( N = transactionCheckMax ) 則 Broker 將丟棄此消息,并在默認情況下同時打印錯誤日志。用戶可以通過重寫 AbstractTransactionCheckListener 類來修改這個行為。

  • 事務消息將在 Broker 配置文件中的參數 transactionMsgTimeout 這樣的特定時間長度之后被檢查。當發送事務消息時,用戶還可以通過設置用戶屬性 CHECK_IMMUNITY_TIME_IN_SECONDS 來改變這個限制,該參數優先于 transactionMsgTimeout 參數。

  • 事務性消息可能不止一次被檢查或消費。

  • 提交給用戶的目標主題消息可能會失敗,目前這依日志的記錄而定。它的高可用性通過 RocketMQ 本身的高可用性機制來保證,如果希望確保事務消息不丟失、并且事務完整性得到保證,建議使用同步的雙重寫入機制。

  • 事務消息的生產者 ID 不能與其他類型消息的生產者 ID 共享。與其他類型的消息不同,事務消息允許反向查詢、MQ服務器能通過它們的生產者 ID 查詢到消費者。

  • 超強干貨來襲 云風專訪:近40年碼齡,通宵達旦的技術人生

    總結

    以上是生活随笔為你收集整理的事务消息的实现的全部內容,希望文章能夠幫你解決所遇到的問題。

    如果覺得生活随笔網站內容還不錯,歡迎將生活随笔推薦給好友。