真实的国产乱ⅩXXX66竹夫人,五月香六月婷婷激情综合,亚洲日本VA一区二区三区,亚洲精品一区二区三区麻豆

成都創(chuàng)新互聯(lián)網(wǎng)站制作重慶分公司

Springboot如何實現(xiàn)RocketMq

小編這次要給大家分享的是Springboot如何實現(xiàn)RocketMq,文章內(nèi)容豐富,感興趣的小伙伴可以來了解一下,希望大家閱讀完這篇文章之后能夠有所收獲。

創(chuàng)新互聯(lián)公司專注為客戶提供全方位的互聯(lián)網(wǎng)綜合服務,包含不限于網(wǎng)站建設、做網(wǎng)站、三臺網(wǎng)絡推廣、小程序設計、三臺網(wǎng)絡營銷、三臺企業(yè)策劃、三臺品牌公關、搜索引擎seo、人物專訪、企業(yè)宣傳片、企業(yè)代運營等,從售前售中售后,我們都將竭誠為您服務,您的肯定,是我們最大的嘉獎;創(chuàng)新互聯(lián)公司為所有大學生創(chuàng)業(yè)者提供三臺建站搭建服務,24小時服務熱線:18980820575,官方網(wǎng)址:www.cdcxhl.com

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

1、pom.xml文件添加依賴

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

    
      org.apache.rocketmq
      rocketmq-remoting
      4.4.0
    

2、yml文件添加rocketmq配置

apache:
 rocketmq:
  #消費者的配置
  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啟動了。。。");
    } 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ā)送響應: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、消費者類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消費者
 * @Date 22:33 2020/5/22
 * @Param
 * @return
 **/
@Component
public class RockerConsumer implements CommandLineRunner {
  /**
   * 消費者
   */
  @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");
      // 程序第一次啟動從消息隊列頭獲取數(shù)據(jù)
      consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
      //可以修改每次消費消息的數(shù)量,默認設置是每次消費一條
      consumer.setConsumeMessageBatchMaxSize(1);

      //在此監(jiān)聽中消費信息,并返回消費的狀態(tài)信息
      consumer.registerMessageListener((MessageListenerConcurrently) (msgs,context)->{
        // 會把不同的消息分別放置到不同的隊列中
        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.測試

請求地址:http://127.0.0.1:8080/rocketMq/myFirstProducer?msg=hello

響應:{"msgId":"C0A8010E1A3818B4AAC2711E8CD50000","sendStatus":"SEND_OK"}

通過rocketMq可視化控制查看:

Springboot如何實現(xiàn)RocketMq

看完這篇關于Springboot如何實現(xiàn)RocketMq的文章,如果覺得文章內(nèi)容寫得不錯的話,可以把它分享出去給更多人看到。


新聞名稱:Springboot如何實現(xiàn)RocketMq
網(wǎng)站鏈接:http://weahome.cn/article/pseipj.html

其他資訊

在線咨詢

微信咨詢

電話咨詢

028-86922220(工作日)

18980820575(7×24)

提交需求

返回頂部