Springboot RocketMq实现过程详解

 更新时间:2020年05月25日 08:41:16   作者:斗战圣猿  
这篇文章主要介绍了Springboot RocketMq实现过程详解,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友可以参考下

首先,在虚拟机上安装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:
  #消费者的配置
  consumer:
   pushConsumer: myConsumer
  #生产者的配置
  producer:
   producerGroup: myGroup
  namesrvAddr: 192.168.233.128:9876

3、生产者类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生产者
 * @Date 22:06 2020/5/22
 * @Param
 * @return
 **/
@Component
public class RocketProducer {
  /**
   * 生产者的组名
   */
  @Value("${apache.rocketmq.producer.producerGroup}")
  private String producerGroup;
  /**
   * NameServer 地址
   */
  @Value("${apache.rocketmq.namesrvAddr}")
  private String namesrvAddr;

  private DefaultMQProducer defaultMQProducer;

  @PostConstruct
  public void defaultMQProducer(){
    //生产者的组名
    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("发送响应:MsgId:" + result.getMsgId() + ",发送状态:" + 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的监听信息,渠道信息
   */
  public void messageListener(){
    DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(pushConsumer);
    consumer.setNamesrvAddr(namesrvAddr);

    try {
      // 订阅PushTopic下Tag为push的消息,都订阅消息
      consumer.subscribe("firstTopic","push");
      // 程序第一次启动从消息队列头获取数据
      consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
      //可以修改每次消费消息的数量,默认设置是每次消费一条
      consumer.setConsumeMessageBatchMaxSize(1);

      //在此监听中消费信息,并返回消费的状态信息
      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中编写发送消息

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如何使用Undertow做服务器

    SpringBoot如何使用Undertow做服务器

    这篇文章主要介绍了SpringBoot如何使用Undertow做服务器,具有很好的参考价值,希望对大家有所帮助。如有错误或未考虑完全的地方,望不吝赐教
    2022-07-07
  • 在非spring环境中调用service中的方法

    在非spring环境中调用service中的方法

    非Spring环境指的是不使用Spring框架来管理和配置应用程序的运行时环境,本文将给大家介绍如何在非spring环境中调用service中的方法,文中有详细实现步骤,需要的朋友可以参考下
    2024-03-03
  • 双token实现token超时策略示例

    双token实现token超时策略示例

    用于restful的app应用无状态无sesion登录示例,需要的朋友可以参考下
    2014-02-02
  • 解决CollectionUtils.isNotEmpty()不存在的问题

    解决CollectionUtils.isNotEmpty()不存在的问题

    这篇文章主要介绍了解决CollectionUtils.isNotEmpty()不存在的问题,具有很好的参考价值,希望对大家有所帮助。如有错误或未考虑完全的地方,望不吝赐教
    2022-02-02
  • Java中Redis的布隆过滤器详解

    Java中Redis的布隆过滤器详解

    这篇文章主要介绍了Java中Redis的布隆过滤器详解,我们经常会把一部分数据放在Redis等缓存,比如产品详情,这样有查询请求进来,我们可以根据产品Id直接去缓存中取数据,而不用读取数据库,这是提升性能最简单,最普遍,也是最有效的做法,需要的朋友可以参考下
    2023-09-09
  • Netty中序列化的作用及自定义协议详解

    Netty中序列化的作用及自定义协议详解

    这篇文章主要介绍了Netty中序列化的作用及自定义协议详解,Netty自身就支持很多种协议比如Http、Websocket等等,但如果用来作为自己的RPC框架通常会自定义协议,所以这也是本文的重点,需要的朋友可以参考下
    2023-12-12
  • java模拟hibernate一级缓存示例分享

    java模拟hibernate一级缓存示例分享

    这篇文章主要介绍了java模拟hibernate一级缓存示例,需要的朋友可以参考下
    2014-03-03
  • 新版本IntelliJ IDEA 构建maven,并用Maven创建一个web项目(图文教程)

    新版本IntelliJ IDEA 构建maven,并用Maven创建一个web项目(图文教程)

    这篇文章主要介绍了新版本IntelliJ IDEA 构建maven,并用Maven创建一个web项目的图文教程,需要的朋友可以参考下
    2018-01-01
  • Spring案例打印机的实现过程详解

    Spring案例打印机的实现过程详解

    这篇文章主要介绍了Spring案例打印机的实现过程详解,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友可以参考下
    2019-10-10
  • 使用fastjson中的JSONPath处理json数据的方法

    使用fastjson中的JSONPath处理json数据的方法

    这篇文章主要介绍了使用fastjson中的JSONPath处理json数据的方法,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学习学习吧
    2020-04-04

最新评论