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

Springboot RocketMq實(shí)現(xiàn)過程詳解

 更新時(shí)間:2020年05月25日 08:41:16   作者:斗戰(zhàn)圣猿  
這篇文章主要介紹了Springboot RocketMq實(shí)現(xiàn)過程詳解,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下

首先,在虛擬機(jī)上安裝rocketmq和rocketMq可視化控制,安裝不做描述。

1、pom.xml文件添加依賴

mq的版本與連接的rocketmq版本保持一致

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

2、yml文件添加rocketmq配置

apache:
 rocketmq:
  #消費(fèi)者的配置
  consumer:
   pushConsumer: myConsumer
  #生產(chǎn)者的配置
  producer:
   producerGroup: myGroup
  namesrvAddr: 192.168.233.128:9876

3、生產(chǎn)者類RocketProducer

package com.zp.springbootdemo.rocketmq;

import com.alibaba.fastjson.JSONObject;
import com.sun.org.apache.xpath.internal.objects.XString;
import org.apache.rocketmq.client.exception.MQBrokerException;
import org.apache.rocketmq.client.exception.MQClientException;
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.apache.rocketmq.remoting.exception.RemotingException;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import org.springframework.util.StopWatch;

import javax.annotation.PostConstruct;
import java.io.UnsupportedEncodingException;

/**
 * @Author zp
 * @Description rocketmq生產(chǎn)者
 * @Date 22:06 2020/5/22
 * @Param
 * @return
 **/
@Component
public class RocketProducer {
  /**
   * 生產(chǎn)者的組名
   */
  @Value("${apache.rocketmq.producer.producerGroup}")
  private String producerGroup;
  /**
   * NameServer 地址
   */
  @Value("${apache.rocketmq.namesrvAddr}")
  private String namesrvAddr;

  private DefaultMQProducer defaultMQProducer;

  @PostConstruct
  public void defaultMQProducer(){
    //生產(chǎn)者的組名
    defaultMQProducer = new DefaultMQProducer(producerGroup);
    defaultMQProducer.setNamesrvAddr(namesrvAddr);
    defaultMQProducer.setVipChannelEnabled(false);
    try {
      defaultMQProducer.start();
      System.out.println("producer啟動(dòng)了。。。");
    } catch (MQClientException e) {
      e.printStackTrace();
    }
  }

  public String send(String topic,String tags,String body) throws UnsupportedEncodingException, InterruptedException, RemotingException, MQClientException, MQBrokerException {
    Message message = new Message(topic,tags,body.getBytes(RemotingHelper.DEFAULT_CHARSET));
    StopWatch stop = new StopWatch();
    stop.start();
    SendResult result = defaultMQProducer.send(message);
    System.out.println("發(fā)送響應(yīng):MsgId:" + result.getMsgId() + ",發(fā)送狀態(tài):" + result.getSendStatus());
    JSONObject jsonObject = new JSONObject();
    jsonObject.put("msgId",result.getMsgId());
    jsonObject.put("sendStatus",result.getSendStatus());
    stop.stop();
    return jsonObject.toJSONString();
  }
}

4、消費(fèi)者類RocketConsumer

package com.zp.springbootdemo.rocketmq;

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.client.exception.MQClientException;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.CommandCustomHeader;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.CommandLineRunner;
import org.springframework.stereotype.Component;

/**
 * @Author zp
 * @Description rocketmq消費(fèi)者
 * @Date 22:33 2020/5/22
 * @Param
 * @return
 **/
@Component
public class RockerConsumer implements CommandLineRunner {
  /**
   * 消費(fèi)者
   */
  @Value("${apache.rocketmq.consumer.pushConsumer}")
  private String pushConsumer; //myConsumer

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

  /**
   * 初始化RocketMq的監(jiān)聽信息,渠道信息
   */
  public void messageListener(){
    DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(pushConsumer);
    consumer.setNamesrvAddr(namesrvAddr);

    try {
      // 訂閱PushTopic下Tag為push的消息,都訂閱消息
      consumer.subscribe("firstTopic","push");
      // 程序第一次啟動(dòng)從消息隊(duì)列頭獲取數(shù)據(jù)
      consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
      //可以修改每次消費(fèi)消息的數(shù)量,默認(rèn)設(shè)置是每次消費(fèi)一條
      consumer.setConsumeMessageBatchMaxSize(1);

      //在此監(jiān)聽中消費(fèi)信息,并返回消費(fèi)的狀態(tài)信息
      consumer.registerMessageListener((MessageListenerConcurrently) (msgs,context)->{
        // 會(huì)把不同的消息分別放置到不同的隊(duì)列中
        for (Message msg:msgs){
          System.out.println("接收到了消息:"+new String(msg.getBody()));

        }
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
      });
      consumer.start();
    } catch (Exception e) {
      e.printStackTrace();
    }
  }
  /**
   * Callback used to run the bean.
   *
   * @param args incoming main method arguments
   * @throws Exception on error
   */
  @Override
  public void run(String... args) throws Exception {
    this.messageListener();
  }
}

5、controller中編寫發(fā)送消息

package com.zp.springbootdemo.rocketmq;

import org.apache.rocketmq.client.exception.MQBrokerException;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.remoting.exception.RemotingException;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;

import java.io.UnsupportedEncodingException;

@RestController
@RequestMapping("/rocketMq")
public class MQController {

  @Autowired
  private RocketProducer producer;

  @RequestMapping("/myFirstProducer")
  public String pushMsg(String msg){
    try {
      System.out.println("======"+msg);
      return producer.send("firstTopic","push",msg);
    } catch (InterruptedException e) {
      e.printStackTrace();
    } catch (RemotingException e) {
      e.printStackTrace();
    } catch (MQClientException e) {
      e.printStackTrace();
    } catch (MQBrokerException e) {
      e.printStackTrace();
    } catch (UnsupportedEncodingException e) {
      e.printStackTrace();
    }
    return "ERROR";
  }
}

6.測(cè)試

請(qǐng)求地址:http://127.0.0.1:8080/rocketMq/myFirstProducer?msg=hello

響應(yīng):{"msgId":"C0A8010E1A3818B4AAC2711E8CD50000","sendStatus":"SEND_OK"}

通過rocketMq可視化控制查看:

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

相關(guān)文章

  • Java Method類及invoke方法原理解析

    Java Method類及invoke方法原理解析

    這篇文章主要介紹了Java Method類及invoke方法原理解析,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2020-08-08
  • 關(guān)于Java項(xiàng)目讀取resources資源文件路徑的那點(diǎn)事

    關(guān)于Java項(xiàng)目讀取resources資源文件路徑的那點(diǎn)事

    這篇文章主要介紹了關(guān)于Java項(xiàng)目讀取resources資源文件路徑的那點(diǎn)事,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-07-07
  • Java 滑動(dòng)窗口最大值的實(shí)現(xiàn)

    Java 滑動(dòng)窗口最大值的實(shí)現(xiàn)

    這篇文章主要介紹了Java 滑動(dòng)窗口最大值,給定一個(gè)數(shù)組 nums,有一個(gè)大小為 k 的滑動(dòng)窗口從數(shù)組的最左側(cè)移動(dòng)到數(shù)組的最右側(cè)。感興趣的可以了解一下
    2021-05-05
  • Mybatis日志模塊的適配器模式詳解

    Mybatis日志模塊的適配器模式詳解

    這篇文章主要介紹了Mybatis日志模塊的適配器模式詳解,,mybatis用了適配器模式來兼容這些框架,適配器模式就是通過組合的方式,將需要適配的類轉(zhuǎn)為使用者能夠使用的接口
    2022-08-08
  • springboot之端口設(shè)置和contextpath的配置方式

    springboot之端口設(shè)置和contextpath的配置方式

    這篇文章主要介紹了springboot之端口設(shè)置和contextpath的配置方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-01-01
  • 詳解Java環(huán)境變量配置方法(Windows)

    詳解Java環(huán)境變量配置方法(Windows)

    這篇文章主要介紹了Java環(huán)境變量配置方法(Windows),文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-03-03
  • Java異常處理中的各種細(xì)節(jié)匯總

    Java異常處理中的各種細(xì)節(jié)匯總

    這篇文章主要給大家介紹了關(guān)于Java異常處理中的各種細(xì)節(jié)的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-01-01
  • Java Http請(qǐng)求傳json數(shù)據(jù)亂碼問題的解決

    Java Http請(qǐng)求傳json數(shù)據(jù)亂碼問題的解決

    這篇文章主要介紹了Java Http請(qǐng)求傳json數(shù)據(jù)亂碼問題的解決,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-09-09
  • 如何基于java向mysql數(shù)據(jù)庫中存取圖片

    如何基于java向mysql數(shù)據(jù)庫中存取圖片

    這篇文章主要介紹了如何基于java向mysql數(shù)據(jù)庫中存取圖片,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2020-02-02
  • SpringCloud feign無法注入接口的問題

    SpringCloud feign無法注入接口的問題

    這篇文章主要介紹了SpringCloud feign無法注入接口的問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-06-06

最新評(píng)論

义乌市| 赞皇县| 枣强县| 旌德县| 会宁县| 延安市| 青龙| 沛县| 咸宁市| 日土县| 千阳县| 叙永县| 大田县| 都安| 桂林市| 绵竹市| 青川县| 太仓市| 会宁县| 长泰县| 嘉峪关市| 卢龙县| 五莲县| 泰宁县| 仪陇县| 兴和县| 周口市| 盈江县| 南京市| 邵阳市| 绿春县| 图们市| 保靖县| 留坝县| 中超| 通化县| 临朐县| 岳普湖县| 昭觉县| 平凉市| 连南|