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

RocketMq事務(wù)消息發(fā)送代碼流程詳解

 更新時(shí)間:2020年07月17日 09:16:56   作者:杯莫停、  
這篇文章主要介紹了RocketMq事務(wù)消息發(fā)送代碼流程詳解,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下

一、RocketMq事務(wù)消息流程:

1、首先會(huì)向broker發(fā)送一個(gè)預(yù)請求消息,消費(fèi)者不可見

2、回調(diào)執(zhí)行本地事務(wù)(比如操作數(shù)據(jù)庫)

3、事務(wù)執(zhí)行成功后,再次發(fā)送消息給broker,告訴broker事務(wù)執(zhí)行成功這個(gè)消息要提交,讓消費(fèi)者可見。如果本地事務(wù)執(zhí)行超時(shí),會(huì)返回一個(gè)unknow,broker會(huì)發(fā)送一個(gè)消息回查,檢查消息是否執(zhí)行成功。

二、RocketMq事務(wù)消息實(shí)例:

1、引入rocketMq相關(guān)的依賴:

<dependency>
  <groupId>org.apache.rocketmq</groupId>
  <artifactId>rocketmq-client</artifactId>
  <version>4.4.0</version>
</dependency>

2、創(chuàng)建一個(gè)TransactionProducer類:

public class TransactionProducer {

  public static void main(String[] args) throws MQClientException, RemotingException, InterruptedException, MQBrokerException, UnsupportedEncodingException {
    //創(chuàng)建生產(chǎn)者并制定組名
    TransactionMQProducer producer = new TransactionMQProducer("rocketMQ_transaction_producer_group");
    //2.指定Nameserver地址
    producer.setNamesrvAddr("192.168.***.***:9876");
    //3、指定消息監(jiān)聽對象用于執(zhí)行本地事務(wù)和消息回查
    TransactionListener listener = new TransactionListenerImol();
    producer.setTransactionListener(listener);
    //4、線程池
    ExecutorService executorService = new ThreadPoolExecutor(2, 5, 100, TimeUnit.SECONDS, new ArrayBlockingQueue<Runnable>(2000), new ThreadFactory() {
      @Override
      public Thread newThread(Runnable r) {
        Thread thread = newThread(r);
        thread.setName("client-tanscation-msg-check-thread");
        return thread;
      }
    });
    producer.setExecutorService(executorService);
    //5、啟動(dòng)producer
    producer.start();

    //6.創(chuàng)建消息對象,指定主題Topic、Tag和消息體 String topic, String tags, String keys, byte[] body
    Message message = new Message("Topic_transaction_demo", //主題
        "Tags", //主要用于消息過濾
        "Key_1", //消息唯一值
        ("hello-transaction").getBytes(RemotingHelper.DEFAULT_CHARSET));

    //7、發(fā)送事務(wù)消息
    TransactionSendResult result = producer.sendMessageInTransaction(message, "hello-transaction");

    producer.shutdown();
  }
}

3、發(fā)送事務(wù)消息還需要一個(gè)事務(wù)監(jiān)聽對象,它實(shí)現(xiàn)TransactionListener 接口,其中有兩個(gè)方法作用分別是執(zhí)行本地事務(wù)和消息回查:

public class TransactionListenerImol implements TransactionListener {
  //存儲(chǔ)事務(wù)狀態(tài)信息 key:事務(wù)id value:當(dāng)前事務(wù)執(zhí)行的狀態(tài)
  private ConcurrentHashMap<String, Integer> localTrans = new ConcurrentHashMap<>();
  //執(zhí)行本地事務(wù)
  @Override
  public LocalTransactionState executeLocalTransaction(Message message, Object o) {
    //事務(wù)id
    String transactionId = message.getTransactionId();
    //0:執(zhí)行中,狀態(tài)未知 1:執(zhí)行成功 2:執(zhí)行失敗
    localTrans.put(transactionId, 0);
    //業(yè)務(wù)執(zhí)行,本地事務(wù),service
    System.out.println("hello-demo-transaction");
    try {
      System.out.println("正在執(zhí)行本地事務(wù)---");
      Thread.sleep(60000*2);
      System.out.println("本地事務(wù)執(zhí)行成功---");
      localTrans.put(transactionId, 1);
    } catch (InterruptedException e) {
      e.printStackTrace();
      localTrans.put(transactionId, 2);
      return LocalTransactionState.ROLLBACK_MESSAGE;
    }
    return LocalTransactionState.COMMIT_MESSAGE;
  }

  //消息回查
  @Override
  public LocalTransactionState checkLocalTransaction(MessageExt messageExt) {
    //獲取對應(yīng)事務(wù)的狀態(tài)信息
    String transactionId = messageExt.getTransactionId();
    //獲取對應(yīng)事務(wù)id執(zhí)行狀態(tài)
    Integer status = localTrans.get(transactionId);
    //消息回查
    System.out.println("消息回查---transactionId:" + transactionId + "狀態(tài):" + status);
    switch (status) {
      case 0:
        return LocalTransactionState.UNKNOW;
      case 1:
        return LocalTransactionState.COMMIT_MESSAGE;
      case 2:
        return LocalTransactionState.ROLLBACK_MESSAGE;
    }
    return LocalTransactionState.UNKNOW;
  }
}

以上就是本文的全部內(nèi)容,希望對大家的學(xué)習(xí)有所幫助,也希望大家多多支持腳本之家。

相關(guān)文章

  • 解決程序啟動(dòng)報(bào)錯(cuò)org.springframework.context.ApplicationContextException: Unable to start web server問題

    解決程序啟動(dòng)報(bào)錯(cuò)org.springframework.context.ApplicationContextExcept

    文章描述了一個(gè)Spring Boot項(xiàng)目在不同環(huán)境下啟動(dòng)時(shí)出現(xiàn)差異的問題,通過分析報(bào)錯(cuò)信息,發(fā)現(xiàn)是由于導(dǎo)入`spring-boot-starter-tomcat`依賴時(shí)定義的scope導(dǎo)致的配置問題,調(diào)整依賴導(dǎo)入配置后,解決了啟動(dòng)錯(cuò)誤
    2024-11-11
  • Java redisson實(shí)現(xiàn)分布式鎖原理詳解

    Java redisson實(shí)現(xiàn)分布式鎖原理詳解

    這篇文章主要介紹了Java redisson實(shí)現(xiàn)分布式鎖原理詳解,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2020-02-02
  • Java實(shí)現(xiàn)多線程的上下文切換

    Java實(shí)現(xiàn)多線程的上下文切換

    這篇文章主要介紹了Java實(shí)現(xiàn)多線程的上下文切換操作,具有很好的參考價(jià)值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2020-09-09
  • Maven?Repository?使用方法

    Maven?Repository?使用方法

    對于Java開發(fā)者來說,Maven?Repository是個(gè)必須掌握的網(wǎng)站,它可以讓開發(fā)者更加方便地管理和維護(hù)?Java?項(xiàng)目的依賴項(xiàng),同時(shí)簡化了項(xiàng)目開發(fā)的過程,這篇文章主要介紹了Maven?Repository?使用方法,需要的朋友可以參考下
    2024-02-02
  • java 中Executor, ExecutorService 和 Executors 間的不同

    java 中Executor, ExecutorService 和 Executors 間的不同

    這篇文章主要介紹了java 中Executor, ExecutorService 和 Executors 間的不同的相關(guān)資料,需要的朋友可以參考下
    2017-06-06
  • IDEA生成servlet程序的實(shí)現(xiàn)步驟

    IDEA生成servlet程序的實(shí)現(xiàn)步驟

    這篇文章主要介紹了IDEA生成servlet程序的實(shí)現(xiàn)步驟,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2021-03-03
  • Java基本語法之內(nèi)部類示例詳解

    Java基本語法之內(nèi)部類示例詳解

    本文帶大家認(rèn)識(shí)Java基本語法——內(nèi)部類,將一個(gè)類定義放在另一類的定義的內(nèi)部,這個(gè)就是內(nèi)部類,內(nèi)部類允許將一些邏輯相關(guān)的類組織在一起,并能夠控制位于內(nèi)部的類的可視性,感興趣的可以了解一下
    2022-03-03
  • 通過JDBC連接oracle數(shù)據(jù)庫的十大技巧

    通過JDBC連接oracle數(shù)據(jù)庫的十大技巧

    通過JDBC連接oracle數(shù)據(jù)庫的十大技巧...
    2006-12-12
  • Java基礎(chǔ)教程之String深度分析

    Java基礎(chǔ)教程之String深度分析

    這篇文章主要給大家介紹了關(guān)于Java基礎(chǔ)教程之String的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-06-06
  • Java Scanner對象中hasNext()與next()方法的使用

    Java Scanner對象中hasNext()與next()方法的使用

    這篇文章主要介紹了Java Scanner對象中hasNext()與next()方法的使用,具有很好的參考價(jià)值,希望對大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-10-10

最新評論

大洼县| 旬阳县| 陈巴尔虎旗| 六安市| 淮滨县| 汝城县| 喀喇沁旗| 清水河县| 宁城县| 潞城市| 临夏县| 特克斯县| 墨江| 平遥县| 肃宁县| 邯郸县| 炉霍县| 济阳县| 扬州市| 阿瓦提县| 陇南市| 南澳县| 孙吴县| 兴安盟| 潞西市| 顺昌县| 鄂托克前旗| 沙湾县| 龙门县| 松江区| 武功县| 三门峡市| 紫云| 婺源县| 南江县| 玉环县| 项城市| 梁河县| 西盟| 永安市| 房产|