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

淺談Springboot整合RocketMQ使用心得

 更新時間:2018年01月15日 16:44:49   作者:HenryZhou2  
本篇文章主要介紹了Springboot整合RocketMQ使用心得,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧

一、阿里云官網(wǎng)---幫助文檔

https://help.aliyun.com/document_detail/29536.html?spm=5176.doc29535.6.555.WWTIUh

按照官網(wǎng)步驟,創(chuàng)建Topic、申請發(fā)布(生產(chǎn)者)、申請訂閱(消費者)

二、代碼

1、配置:

public class MqConfig {
  /**
   * 啟動測試之前請?zhí)鎿Q如下 XXX 為您的配置
   */
  public static final String PUBLIC_TOPIC = "test";//公網(wǎng)測試
  public static final String PUBLIC_PRODUCER_ID = "PID_SCHEDULER";
  public static final String PUBLIC_CONSUMER_ID = "CID_SERVICE";

  public static final String ACCESS_KEY = "123";
  public static final String SECRET_KEY = "123";
  public static final String TAG = "";
  public static final String THREAD_NUM = "25";//消費端線程數(shù)
  /**
   * ONSADDR 請根據(jù)不同Region進行配置
   * 公網(wǎng)測試: http://onsaddr-internet.aliyun.com/rocketmq/nsaddr4client-internet
   * 公有云生產(chǎn): http://onsaddr-internal.aliyun.com:8080/rocketmq/nsaddr4client-internal
   * 杭州金融云: http://jbponsaddr-internal.aliyun.com:8080/rocketmq/nsaddr4client-internal
   * 深圳金融云: http://mq4finance-sz.addr.aliyun.com:8080/rocketmq/nsaddr4client-internal
   */
  public static final String ONSADDR = "http://onsaddr-internal.aliyun.com:8080/rocketmq/nsaddr4client-internal";
}

ONSADDR 阿里云用 公有云生產(chǎn),測試用公網(wǎng)

不同的業(yè)務(wù)可以設(shè)置不同的tag,但是如果發(fā)送消息量大的話,建議新建TOPIC

2、生產(chǎn)者

方式1:

配置文件:producer.xml

<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE beans PUBLIC "-//SPRING//DTD BEAN//EN" "http://www.springframework.org/dtd/spring-beans.dtd">
<beans>
  <bean id="producer" class="com.aliyun.openservices.ons.api.bean.ProducerBean"
     init-method="start" destroy-method="shutdown">
    <property name="properties">
      <map>
        <entry key="ProducerId" value="" /> <!-- PID,請?zhí)鎿Q -->
        <entry key="AccessKey" value="" /> <!-- ACCESS_KEY,請?zhí)鎿Q -->
        <entry key="SecretKey" value="" /> <!-- SECRET_KEY,請?zhí)鎿Q -->
        <!--PropertyKeyConst.ONSAddr 請根據(jù)不同Region進行配置
         公網(wǎng)測試: http://onsaddr-internet.aliyun.com/rocketmq/nsaddr4client-internet
         公有云生產(chǎn): http://onsaddr-internal.aliyun.com:8080/rocketmq/nsaddr4client-internal
         杭州金融云: http://jbponsaddr-internal.aliyun.com:8080/rocketmq/nsaddr4client-internal
         深圳金融云: http://mq4finance-sz.addr.aliyun.com:8080/rocketmq/nsaddr4client-internal -->
        <entry key="ONSAddr" value="http://onsaddr-internal.aliyun.com:8080/rocketmq/nsaddr4client-internal"/>
      </map>
    </property>
  </bean>
</beans>

啟動方式1,在使用類的全局里設(shè)置:

//初始化生產(chǎn)者
  private ApplicationContext ctx;
  private ProducerBean producer;

  @Value("${producerConfig.enabled}")//開關(guān),spring配置項,true為開啟,false關(guān)閉
  private boolean producerConfigEnabled;

  @PostConstruct
  public void init(){
    if (true == producerConfigEnabled) {
      ctx = new ClassPathXmlApplicationContext("producer.xml");
      producer = (ProducerBean) ctx.getBean("producer");
    }
  }

PS:最近發(fā)現(xiàn)一個坑,如果producer用上面這種方式啟動的話,一旦啟動的多了,會造成fullGC,所以可以換成下面這種注解方式啟動,在用到的地方手動start、shutdown

方式2:配置類(不需要xml)

@Configuration
public class ProducerBeanConfig {

  @Value("${openservices.ons.producerBean.producerId}")
  private String producerId;

  @Value("${openservices.ons.producerBean.accessKey}")
  private String accessKey;

  @Value("${openservices.ons.producerBean.secretKey}")
  private String secretKey;

  private ProducerBean producerBean;

  @Value("${openservices.ons.producerBean.ONSAddr}")
  private String ONSAddr;

  @Bean
  public ProducerBean oneProducer() {
    ProducerBean producerBean = new ProducerBean();
    Properties properties = new Properties();
    properties.setProperty(PropertyKeyConst.ProducerId, producerId);
    properties.setProperty(PropertyKeyConst.AccessKey, accessKey);
    properties.setProperty(PropertyKeyConst.SecretKey, secretKey);
    properties.setProperty(PropertyKeyConst.ONSAddr, ONSAddr);

    producerBean.setProperties(properties);
    return producerBean;
  }
}

PS:經(jīng)過這次雙11發(fā)現(xiàn),以上2種方式在大數(shù)據(jù)量,多線程情況下都不太適用, 性能很差,所以推薦用3

方式3:(不需要xml)

@Component
public class ProducerBeanSingleTon {

  @Value("${openservices.ons.producerBean.producerId}")
  private String producerId;

  @Value("${openservices.ons.producerBean.accessKey}")
  private String accessKey;

  @Value("${openservices.ons.producerBean.secretKey}")
  private String secretKey;

  @Value("${openservices.ons.producerBean.ONSAddr}")
  private String ONSAddr;

  private static Producer producer;

  private static class SingletonHolder {
    private static final ProducerBeanSingleTon INSTANCE = new ProducerBeanSingleTon();
  }

  private ProducerBeanSingleTon (){}

  public static final ProducerBeanSingleTon getInstance() {
    return SingletonHolder.INSTANCE;
  }

  @PostConstruct
  public void init(){
    // producer 實例配置初始化
    Properties properties = new Properties();
    //您在控制臺創(chuàng)建的Producer ID
    properties.setProperty(PropertyKeyConst.ProducerId, producerId);
    // AccessKey 阿里云身份驗證,在阿里云服務(wù)器管理控制臺創(chuàng)建
    properties.setProperty(PropertyKeyConst.AccessKey, accessKey);
    // SecretKey 阿里云身份驗證,在阿里云服務(wù)器管理控制臺創(chuàng)建
    properties.setProperty(PropertyKeyConst.SecretKey, secretKey);
    //設(shè)置發(fā)送超時時間,單位毫秒
    properties.setProperty(PropertyKeyConst.SendMsgTimeoutMillis, "3000");
    // 設(shè)置 TCP 接入域名(此處以公共云生產(chǎn)環(huán)境為例)
    properties.setProperty(PropertyKeyConst.ONSAddr, ONSAddr);
    producer = ONSFactory.createProducer(properties);
    // 在發(fā)送消息前,必須調(diào)用start方法來啟動Producer,只需調(diào)用一次即可
    producer.start();
  }

  public Producer getProducer(){
    return producer;
  }
}

spring配置

spring.jpa.properties.hibernate.dialect = org.hibernate.dialect.MySQL5Dialect

consumerConfig.enabled = true

producerConfig.enabled = true #方式1:

scheduling.enabled = false

#方式2、3:rocketMQ \u516C\u7F51\u914D\u7F6E
openservices.ons.producerBean.producerId = pid
openservices.ons.producerBean.accessKey = 
openservices.ons.producerBean.secretKey = 

openservices.ons.producerBean.ONSAddr = 公網(wǎng)、杭州公有云生產(chǎn)

方式1投遞消息代碼:

 try {
   String jsonC = JsonUtils.toJson(elevenMessage);
   Message message = new Message(MqConfig.TOPIC, MqConfig.TAG, jsonC.getBytes());
   SendResult sendResult = producer.send(message);
   if (sendResult != null) {
     logger.info(".Send mq message success!”;

   } else {
     logger.warn(".sendResult is null.........");
   }
   } catch (Exception e) {
      logger.warn("DoubleElevenAllPreService");
      Thread.sleep(1000);//如果有異常,休眠1秒
   }

方式2投遞消息代碼:(可以每發(fā)1000個啟動/關(guān)閉一次)

   producerBean.start();
try {
   String jsonC = JsonUtils.toJson(elevenMessage);
   Message message = new Message(MqConfig.TOPIC, MqConfig.TAG, jsonC.getBytes());
   SendResult sendResult = producer.send(message);
   if (sendResult != null) {
     logger.info(".Send mq message success!”;

   } else {
     logger.warn(".sendResult is null.........");
   }
   } catch (Exception e) {
      logger.warn("DoubleElevenAllPreService");
      Thread.sleep(1000);//如果有異常,休眠1秒
   }

   producerBean.shutdown();

方式3:投遞消息

 try {
   String jsonC = JsonUtils.toJson(elevenMessage);
   Message message = new Message(MqConfig.TOPIC, MqConfig.TAG, jsonC.getBytes());
   Producer producer = ProducerBeanSingleTon.getInstance().getProducer();
   SendResult sendResult = producer.send(message);
   if (sendResult != null) {
     logger.info("DoubleElevenMidService.Send mq message success! Topic is:"”;

   } else {
     logger.warn("DoubleElevenMidService.sendResult is null.........");
   }
   } catch (Exception e) {
     logger.error("DoubleElevenMidService Thread.sleep 1 s___error is "+e.getMessage(), e);
     Thread.sleep(1000);//如果有異常,休眠1秒
   }

發(fā)送消息的代碼一定要捕獲異常,不然會重復(fù)發(fā)送。

這里的TOPIC用自己創(chuàng)建的,elevenMessage是要發(fā)送的內(nèi)容,我這里是自己建的對象

3、消費者

配置啟動類:

@Configuration
@ConditionalOnProperty(value = "consumerConfig.enabled", havingValue = "true", matchIfMissing = true)
public class ConsumerConfig {

  private Logger logger = LoggerFactory.getLogger(LoggerAppenderType.smsdist.name());

  @Bean
  public Consumer consumerFactory(){//不同消費者 這里不能重名
    Properties consumerProperties = new Properties();
    consumerProperties.setProperty(PropertyKeyConst.ConsumerId, MqConfig.CONSUMER_ID);
    consumerProperties.setProperty(PropertyKeyConst.AccessKey, MqConfig.ACCESS_KEY);
    consumerProperties.setProperty(PropertyKeyConst.SecretKey, MqConfig.SECRET_KEY);
    //consumerProperties.setProperty(PropertyKeyConst.ConsumeThreadNums,MqConfig.THREAD_NUM);
    consumerProperties.setProperty(PropertyKeyConst.ONSAddr, MqConfig.ONSADDR);
    Consumer consumer = ONSFactory.createConsumer(consumerProperties);
    consumer.subscribe(MqConfig.TOPIC, MqConfig.TAG, new DoubleElevenMessageListener());//new對應(yīng)的監(jiān)聽器
    consumer.start();
    logger.info("ConsumerConfig start success.");
    

    return consumer;

  }
}

CID和ONSADDR一點要選對,用自己的,消費者線程數(shù)等可以在這里配置

創(chuàng)建消息監(jiān)聽器類,消費消息:

@Component
public class MessageListener implements MessageListener {
  private Logger logger = LoggerFactory.getLogger("remind");

  protected static ElevenReposity elevenReposity;
  @Resource
  public void setElevenReposity(ElevenReposity elevenReposity){
    MessageListener .elevenReposity=elevenReposity;
  }


  @Override
  public Action consume(Message message, ConsumeContext consumeContext) {

    if(message.getTopic().equals("自己的TOPIC")){//避免消費到其他消息 json轉(zhuǎn)換報錯
      try {

      byte[] body = message.getBody();
      String res = new String(body);
      
      //res 是生產(chǎn)者傳過來的消息內(nèi)容

        //業(yè)務(wù)代碼

      }else{
        logger.warn("!");
      }

      } catch (Exception e) {
        logger.error("MessageListener.consume error:" + e.getMessage(), e);
      }

      logger.info("MessageListener.Receive message”);
      //如果想測試消息重投的功能,可以將Action.CommitMessage 替換成Action.ReconsumeLater
      return Action.CommitMessage;
    }else{
      logger.warn();
      return Action.ReconsumeLater;
    }

  }

注意,由于消費者是多線程的,所以對象要用static+set注入,把對象的級別提升到進程,這樣多個線程就可以共用,但是無法調(diào)用父類的方法和變量

消費者狀態(tài)可以查看消費者是否連接成功,消費是否延遲,消費速度等

重置消費位點可以清空所有消息

三、注意事項

1、發(fā)送的消息體 最大為256KB

2、消息最多存在3天

3、消費端默認線程數(shù)是20

4、如果運行過程中出現(xiàn)java掛掉或者cpu占用異常高,可以在發(fā)送消息的時候,每發(fā)送1000條讓線程休息1s

5、本地測試或啟動的時候,把ONSADDR換成公網(wǎng),不然報錯無法啟動

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

相關(guān)文章

  • Java查找不重復(fù)無序數(shù)組中是否存在兩個數(shù)字的和為某個值

    Java查找不重復(fù)無序數(shù)組中是否存在兩個數(shù)字的和為某個值

    今天小編就為大家分享一篇關(guān)于Java查找不重復(fù)無序數(shù)組中是否存在兩個數(shù)字的和為某個值,小編覺得內(nèi)容挺不錯的,現(xiàn)在分享給大家,具有很好的參考價值,需要的朋友一起跟隨小編來看看吧
    2019-01-01
  • Spring Security實現(xiàn)自動登陸功能示例

    Spring Security實現(xiàn)自動登陸功能示例

    自動登錄在很多網(wǎng)站和APP上都能用的到,解決了用戶每次輸入賬號密碼的麻煩。本文就使用Spring Security實現(xiàn)自動登陸功能,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2021-11-11
  • SpringBoot依賴管理的源碼解析

    SpringBoot依賴管理的源碼解析

    這篇文章主要介紹了SpringBoot依賴管理的源碼解析,maven提供了一套依賴管理機制,通過在pom.xml定義坐標,通過坐標從互聯(lián)網(wǎng)的中央倉庫下載依賴的構(gòu)件(jar包),規(guī)范去管理依賴所有構(gòu)件,這就叫依賴管理,需要的朋友可以參考下
    2023-04-04
  • 一文講解Java的String、StringBuffer和StringBuilder的使用與區(qū)別

    一文講解Java的String、StringBuffer和StringBuilder的使用與區(qū)別

    String是不可變的字符序列,而StringBuffer和StringBuilder是可變的字符序列,本文就來詳細的介紹一下Java的String、StringBuffer和StringBuilder的使用與區(qū)別,感興趣的可以了解一下
    2024-03-03
  • Java實現(xiàn)貪吃蛇大作戰(zhàn)小游戲(附源碼)

    Java實現(xiàn)貪吃蛇大作戰(zhàn)小游戲(附源碼)

    今天給大家?guī)淼氖切№椖渴?nbsp;基于Java+Swing+IO流實現(xiàn) 的貪吃蛇大作戰(zhàn)小游戲。實現(xiàn)了界面可視化、基本的吃食物功能、死亡功能、移動功能、積分功能,并額外實現(xiàn)了主動加速和鼓勵機制,需要的可以參考一下
    2022-07-07
  • SpringDataJpa多表操作的實現(xiàn)

    SpringDataJpa多表操作的實現(xiàn)

    開發(fā)過程中會有很多多表的操作,他們之間有著各種關(guān)系,本文主要介紹了SpringDataJpa多表操作的實現(xiàn),具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2021-11-11
  • Java中API的使用方法詳情

    Java中API的使用方法詳情

    這篇文章主要介紹了Java中API的使用方法詳情,指的就是?JDK?中提供的各種功能的?Java類,這些類將底層的實現(xiàn)封裝了起來,我們不需要關(guān)心這些類是如何實現(xiàn)的,只需要學習這些類如何使用即可,我們可以通過幫助文檔來學習這些API如何使用,需要的朋友可以參考下
    2022-04-04
  • Java利用反射獲取object的屬性和值代碼示例

    Java利用反射獲取object的屬性和值代碼示例

    這篇文章主要介紹了Java利用反射獲取object的屬性和值代碼示例,具有一定借鑒價值,需要的朋友可以參考下。
    2017-12-12
  • SpringBoot之自定義Filter獲取請求參數(shù)與響應(yīng)結(jié)果案例詳解

    SpringBoot之自定義Filter獲取請求參數(shù)與響應(yīng)結(jié)果案例詳解

    這篇文章主要介紹了SpringBoot之自定義Filter獲取請求參數(shù)與響應(yīng)結(jié)果案例詳解,本篇文章通過簡要的案例,講解了該項技術(shù)的了解與使用,以下就是詳細內(nèi)容,需要的朋友可以參考下
    2021-09-09
  • Java無界阻塞隊列DelayQueue詳細解析

    Java無界阻塞隊列DelayQueue詳細解析

    這篇文章主要介紹了Java無界阻塞隊列DelayQueue詳細解析,DelayQueue是一個支持時延獲取元素的無界阻塞隊列,隊列使用PriorityQueue來實現(xiàn),隊列中的元素必須實現(xiàn)Delayed接口,在創(chuàng)建元素時可以指定多久才能從隊列中獲取當前元素,需要的朋友可以參考下
    2023-12-12

最新評論

定日县| 太湖县| 汶川县| 辽阳市| 龙陵县| 高邑县| 印江| 仁化县| 永登县| 翁源县| 日喀则市| 交城县| 怀远县| 南宁市| 丽江市| 四子王旗| 呼图壁县| 汽车| 娱乐| 绵阳市| 扎鲁特旗| 武川县| 沁阳市| 土默特左旗| 东明县| 韶山市| 格尔木市| 沅江市| 临洮县| 安龙县| 宁海县| 东莞市| 溧阳市| 牟定县| 综艺| 南通市| 漠河县| 三亚市| 静安区| 全椒县| 明星|