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

springBoot整合RocketMQ及坑的示例代碼

 更新時(shí)間:2018年11月12日 15:02:38   作者:龍俊潔  
這篇文章主要介紹了springBoot整合RocketMQ及坑的示例代碼,小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過來看看吧

版本:

  • JDK:1.8
  • springBoot:1.5.10
  • rocketMQ:4.2.0

pom 配置:    

<parent>
 <groupId>org.springframework.boot</groupId>
 <artifactId>spring-boot-starter-parent</artifactId>
 <version>1.5.10.RELEASE</version>
</parent>
<dependency>
  <groupId>org.apache.rocketmq</groupId>
  <artifactId>rocketmq-client</artifactId>
  <version>4.2.0</version>
</dependency>

application.properties  配置:

# 消費(fèi)者的組名
apache.rocketmq.consumer.PushConsumer=PushConsumer
# 生產(chǎn)者的組名
apache.rocketmq.producer.producerGroup=Producer
# NameServer地址
apache.rocketmq.namesrvAddr=localhost:9876

java代碼:

生產(chǎn)者

package test.config.rocketmq;

import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import org.springframework.util.StopWatch;
import javax.annotation.PostConstruct;

@Component
public class RocketMQClient {
  /**
   * 生產(chǎn)者的組名
   */
  @Value("${apache.rocketmq.producer.producerGroup}")
  private String producerGroup;

  /**
   * NameServer 地址
   */
  @Value("${apache.rocketmq.namesrvAddr}")
  private String namesrvAddr;

  @PostConstruct
  public void defaultMQProducer() {
    //生產(chǎn)者的組名
    DefaultMQProducer producer = new DefaultMQProducer(producerGroup);
    //指定NameServer地址,多個(gè)地址以 ; 隔開
    producer.setNamesrvAddr(namesrvAddr);
    producer.setVipChannelEnabled(false);
    try {
      /**
       * Producer對(duì)象在使用之前必須要調(diào)用start初始化,初始化一次即可
       * 注意:切記不可以在每次發(fā)送消息時(shí),都調(diào)用start方法
       */
      producer.start();

      //創(chuàng)建一個(gè)消息實(shí)例,包含 topic、tag 和 消息體
      //如下:topic 為 "TopicTest",tag 為 "push"
      Message message = new Message("TopicTest", "push", "發(fā)送消息----zhisheng-----".getBytes(RemotingHelper.DEFAULT_CHARSET));

      StopWatch stop = new StopWatch();
      stop.start();

      for (int i = 0; i < 1; i++) {
        SendResult result = producer.send(message);
        System.out.println("發(fā)送響應(yīng):MsgId:" + result.getMsgId() + ",發(fā)送狀態(tài):" + result.getSendStatus());
      }
      stop.stop();
      System.out.println("----------------發(fā)送一萬(wàn)條消息耗時(shí):" + stop.getTotalTimeMillis());
    } catch (Exception e) {
      e.printStackTrace();
    } finally {
      producer.shutdown();
    }
  }
}

消費(fèi)者: 

import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;

import javax.annotation.PostConstruct;


@Component
public class RocketMQServer {
  /**
   * 消費(fèi)者的組名
   */
  @Value("${apache.rocketmq.consumer.PushConsumer}")
  private String consumerGroup;

  /**
   * NameServer 地址
   */
  @Value("${apache.rocketmq.namesrvAddr}")
  private String namesrvAddr;

  @PostConstruct
  public void defaultMQPushConsumer() {
    //消費(fèi)者的組名
    DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(consumerGroup);

    //指定NameServer地址,多個(gè)地址以 ; 隔開
    consumer.setNamesrvAddr(namesrvAddr);
    consumer.setVipChannelEnabled(false);
    try {
      //訂閱PushTopic下Tag為push的消息
      consumer.subscribe("TopicTest", "push");

      //設(shè)置Consumer第一次啟動(dòng)是從隊(duì)列頭部開始消費(fèi)還是隊(duì)列尾部開始消費(fèi)
      //如果非第一次啟動(dòng),那么按照上次消費(fèi)的位置繼續(xù)消費(fèi)
      consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
      consumer.registerMessageListener((MessageListenerConcurrently) (list, context) -> {
        try {
          for (MessageExt messageExt : list) {

            System.out.println("messageExt: " + messageExt);//輸出消息內(nèi)容

            String messageBody = new String(messageExt.getBody(), RemotingHelper.DEFAULT_CHARSET);

            System.out.println("消費(fèi)響應(yīng):msgId : " + messageExt.getMsgId() + ", msgBody : " + messageBody);//輸出消息內(nèi)容
          }
        } catch (Exception e) {
          e.printStackTrace();
          return ConsumeConcurrentlyStatus.RECONSUME_LATER; //稍后再試
        }
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; //消費(fèi)成功
      });
      consumer.start();
    } catch (Exception e) {
      e.printStackTrace();
    }
  }
}

掉坑總結(jié):

1.rocketMQ啟動(dòng)時(shí),命令不是  mqbroker -n 127.0.0.1:9876

         正確應(yīng)該是:mqbroker -n 127.0.0.1:9876 butiautoCreateTopicEnable=true

         否則會(huì)拋出:No route info of this topic, TopicTest

2.客戶端連接時(shí)拋出異常

        org.apache.rocketmq.client.exception.MQClientException: 

        Send [3] times, still failed, cost [3180]ms, Topic: TopicTest, BrokersSent: \

        [WIN-93CGO0S5G25, WIN-93CGO0S5G25, WIN-93CGO0S5G25]

解決方式兩種

1.producer.setVipChannelEnabled(false); 生產(chǎn)者和消費(fèi)者添加這行代買。

2.降rocketmq版本,降成3.2.6

關(guān)于spring.rocketmq.name-server的坑

看下圖:

注意:

如果你是SpringBoot2.0+的框架,或者是JDK10。

你需要將你自己的項(xiàng)目配置文件中的,spring.rocketmq.name-server改成

spring.rocketmq.nameServer。注意是nameServer。

不然就會(huì)報(bào)各種稀奇古怪的bug。

關(guān)于啟動(dòng)報(bào)內(nèi)存不足的錯(cuò)

在安裝啟動(dòng)Name Server和Broker的時(shí)候,一定要修改配置文件,不然內(nèi)存會(huì)爆炸。

Native memory allocation (mmap) failed to map 8589934592 bytes for committing reserved memory 

將下面的配置文件根據(jù)你的需要改

我這里以前默認(rèn)是Xms4g,都是g,我修改到m就行了。

JAVA_OPT="${JAVA_OPT} -server -Xms256m -Xmx256m -Xmn128m -XX:MetaspaceSize=128m -XX:MaxMetaspaceSize=320m"

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

相關(guān)文章

  • java web實(shí)現(xiàn)用戶權(quán)限管理

    java web實(shí)現(xiàn)用戶權(quán)限管理

    這篇文章主要介紹了java web實(shí)現(xiàn)用戶權(quán)限管理,設(shè)計(jì)并實(shí)現(xiàn)一套簡(jiǎn)單的權(quán)限管理功能,感興趣的小伙伴們可以參考一下
    2015-11-11
  • java多線程Future和Callable類示例分享

    java多線程Future和Callable類示例分享

    JAVA多線程實(shí)現(xiàn)方式主要有三種:繼承Thread類、實(shí)現(xiàn)Runnable接口、使用ExecutorService、Callable、Future實(shí)現(xiàn)有返回結(jié)果的多線程。其中前兩種方式線程執(zhí)行完后都沒有返回值,只有最后一種是帶返回值的。今天我們就來研究下Future和Callable的實(shí)現(xiàn)方法
    2016-01-01
  • java9中g(shù)c log參數(shù)遷移

    java9中g(shù)c log參數(shù)遷移

    本篇文章給大家詳細(xì)講述了java9中g(shù)c log參數(shù)遷移的相關(guān)知識(shí)點(diǎn),對(duì)此有需要的朋友可以參考學(xué)習(xí)下。
    2018-03-03
  • Java中線程上下文類加載器超詳細(xì)講解使用

    Java中線程上下文類加載器超詳細(xì)講解使用

    這篇文章主要介紹了Java中線程上下文類加載器,類加載器負(fù)責(zé)讀取Java字節(jié)代碼,并轉(zhuǎn)換成java.lang.Class類的一個(gè)實(shí)例的代碼模塊。本文主要和大家聊聊JVM類加載器ClassLoader的使用,需要的可以了解一下
    2022-12-12
  • pdf2swf+flexpapers實(shí)現(xiàn)類似百度文庫(kù)pdf在線閱讀

    pdf2swf+flexpapers實(shí)現(xiàn)類似百度文庫(kù)pdf在線閱讀

    這篇文章主要介紹了pdf2swf+flexpapers實(shí)現(xiàn)類似百度文庫(kù)pdf在線閱讀的相關(guān)資料,需要的朋友可以參考下
    2014-10-10
  • @GrpcServise?注解的作用和使用示例詳解

    @GrpcServise?注解的作用和使用示例詳解

    @GrpcService 是一個(gè) Spring Boot 處理器,它會(huì)查找實(shí)現(xiàn)了 grpc::BindableService 接口的類,并將其包裝成一個(gè) Spring Bean 對(duì)象,這篇文章主要介紹了@GrpcServise?注解的作用和使用,需要的朋友可以參考下
    2023-05-05
  • 配置SpringBoot中的Jackson序列化方式

    配置SpringBoot中的Jackson序列化方式

    本文介紹了在SpringBoot開發(fā)中,如何通過Jackson自定義JSON序列化和反序列化配置,解決常見問題,如日期格式、空值處理、數(shù)據(jù)精度等,總結(jié) Jackson默認(rèn)行為通常滿足需求,但通過配置ObjectMapper和JsonConstants,可以靈活調(diào)整序列化邏輯,滿足具體應(yīng)用需求
    2025-09-09
  • Java設(shè)計(jì)模式之裝飾模式詳解

    Java設(shè)計(jì)模式之裝飾模式詳解

    這篇文章主要介紹了Java設(shè)計(jì)模式中的裝飾者模式,裝飾者模式即Decorator Pattern,裝飾模式是在不必改變?cè)愇募褪褂美^承的情況下,動(dòng)態(tài)地?cái)U(kuò)展一個(gè)對(duì)象的功能,裝飾模式又名包裝模式。裝飾器模式以對(duì)客戶端透明的方式拓展對(duì)象的功能,是繼承關(guān)系的一種替代方案
    2022-08-08
  • Java JUC中操作List安全類的集合案例

    Java JUC中操作List安全類的集合案例

    這篇文章主要介紹了JUC中操作List安全類的集合案例,本文羅列了不安全的集合和安全的集合進(jìn)行對(duì)比,以及Java中提供的安全措施,需要的朋友可以參考下
    2021-07-07
  • Java求1+2!+3!+...+20!的和的代碼

    Java求1+2!+3!+...+20!的和的代碼

    這篇文章主要介紹了Java求1+2!+3!+...+20!的和的代碼,需要的朋友可以參考下
    2017-02-02

最新評(píng)論

临沂市| 蕉岭县| 黄石市| 会东县| 呼玛县| 临桂县| 金昌市| 谢通门县| 江源县| 泌阳县| 加查县| 固原市| 新乐市| 鄂州市| 张家界市| 万州区| 安西县| 陕西省| 泸州市| 稻城县| 余庆县| 菏泽市| 微山县| 双鸭山市| 大连市| 长治县| 桓仁| 临夏市| 开阳县| 微山县| 牙克石市| 类乌齐县| 东乡族自治县| 晋宁县| 南京市| 安国市| 边坝县| 灌阳县| 宝坻区| 大英县| 北辰区|