springboot-rocketmq整合

1、application.properties
spring.application.name = demoTest

mybatis

spring.profiles.active=dev
spring.datasource.url=jdbc:mysql://localhost:3306/saasboard
spring.datasource.username=root
spring.datasource.password=root
spring.datasource.driver-class-name=com.mysql.jdbc.Driver
spring.datasource.type=com.alibaba.druid.pool.DruidDataSource

region redis

spring.redis.host=192.168.207.18
spring.redis.port=6379

spring.redis.password=group@123

spring.redis.database=0
spring.redis.pool.max-active=150
spring.redis.pool.max-idle=30
spring.redis.pool.max-wait=3000
spring.redis.pool.min-idle=10

producer

该应用是否启用生产者

rocketmq.producer.isOnOff=on

发送同一类消息的设置为同一个group,保证唯一,默认不需要设置,rocketmq会使用ip@pid(pid代表jvm名字)作为唯一标示

rocketmq.producer.groupName=${spring.application.name}

mq的nameserver地址

rocketmq.producer.namesrvAddr=192.168.205.196:9876

消息最大长度 默认1024*4(4M)

rocketmq.producer.maxMessageSize=4096

发送消息超时时间,默认3000

rocketmq.producer.sendMsgTimeout=3000

发送消息失败重试次数,默认2

rocketmq.producer.retryTimesWhenSendFailed=2

consumer

该应用是否启用消费者

rocketmq.consumer.isOnOff=on
rocketmq.consumer.groupName=${spring.application.name}

mq的nameserver地址

rocketmq.consumer.namesrvAddr=192.168.205.196:9876

该消费者订阅的主题和tags(""号表示订阅该主题下所有的tags),格式:topictag1||tag2||tag3;topic2;

rocketmq.consumer.topics=DemoTopic~*;
rocketmq.consumer.consumeThreadMin=20
rocketmq.consumer.consumeThreadMax=64

设置一次消费消息的条数,默认为1条

rocketmq.consumer.consumeMessageBatchMaxSize=1

2、pom.xml文件
<dependency>
<groupId>com.alibaba.rocketmq</groupId>
<artifactId>rocketmq-client</artifactId>
<version>3.2.6</version>
</dependency>
3、消费端配置
package com.example.demo.rocketmq.consumer;

import com.alibaba.rocketmq.client.consumer.DefaultMQPushConsumer;
import com.alibaba.rocketmq.client.exception.MQClientException;
import com.alibaba.rocketmq.common.consumer.ConsumeFromWhere;
import com.example.demo.rocketmq.constants.RocketMQErrorEnum;
import com.example.demo.rocketmq.consumer.processor.MQConsumeMsgListenerProcessor;
import com.example.demo.rocketmq.exception.RocketMQException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.SpringBootConfiguration;
import org.springframework.context.annotation.Bean;
import org.springframework.util.StringUtils;

/**

  • 消费者Bean配置

  • .

  • Copyright: Copyright (c) 2017 zteits

  • @ClassName: MQConsumerConfiguration

  • @Description:

  • @version: v1.0.0

  • @author: zhaowg

  • @date: 2018年3月2日 下午11:48:32

  • Modification History:

  • Date Author Version Description
    ---------------------------------------------------------

  • 2018年3月2日 zhaowg v1.0.0 创建
    */
    @SpringBootConfiguration
    public class MQConsumerConfiguration {
    public static final Logger LOGGER = LoggerFactory.getLogger(MQConsumerConfiguration.class);
    @Value("{rocketmq.consumer.namesrvAddr}") private String namesrvAddr; @Value("{rocketmq.consumer.groupName}")
    private String groupName;
    @Value("{rocketmq.consumer.consumeThreadMin}") private int consumeThreadMin; @Value("{rocketmq.consumer.consumeThreadMax}")
    private int consumeThreadMax;
    @Value("{rocketmq.consumer.topics}") private String topics; @Value("{rocketmq.consumer.consumeMessageBatchMaxSize}")
    private int consumeMessageBatchMaxSize;
    @Autowired
    private MQConsumeMsgListenerProcessor mqMessageListenerProcessor;

    @Bean
    public DefaultMQPushConsumer getRocketMQConsumer() throws RocketMQException {
    if (StringUtils.isEmpty(groupName)){
    throw new RocketMQException(RocketMQErrorEnum.PARAMM_NULL,"groupName is null !!!",false);
    }
    if (StringUtils.isEmpty(namesrvAddr)){
    throw new RocketMQException(RocketMQErrorEnum.PARAMM_NULL,"namesrvAddr is null !!!",false);
    }
    if(StringUtils.isEmpty(topics)){
    throw new RocketMQException(RocketMQErrorEnum.PARAMM_NULL,"topics is null !!!",false);
    }
    DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(groupName);
    consumer.setNamesrvAddr(namesrvAddr);
    consumer.setConsumeThreadMin(consumeThreadMin);
    consumer.setConsumeThreadMax(consumeThreadMax);
    consumer.registerMessageListener(mqMessageListenerProcessor);
    /**
    * 设置Consumer第一次启动是从队列头部开始消费还是队列尾部开始消费
    * 如果非第一次启动,那么按照上次消费的位置继续消费
    /
    consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET);
    /
    *
    * 设置消费模型,集群还是广播,默认为集群
    /
    //consumer.setMessageModel(MessageModel.CLUSTERING);
    /
    *
    * 设置一次消费消息的条数,默认为1条
    /
    consumer.setConsumeMessageBatchMaxSize(consumeMessageBatchMaxSize);
    try {
    /
    *
    * 设置该消费者订阅的主题和tag,如果是订阅该主题下的所有tag,则tag使用*;如果需要指定订阅该主题下的某些tag,则使用||分割,例如tag1||tag2||tag3
    */
    String[] topicTagsArr = topics.split(";");
    for (String topicTags : topicTagsArr) {
    String[] topicTag = topicTags.split("~");
    consumer.subscribe(topicTag[0],topicTag[1]);
    }
    consumer.start();
    LOGGER.info("consumer is start !!! groupName:{},topics:{},namesrvAddr:{}",groupName,topics,namesrvAddr);
    }catch (MQClientException e){
    LOGGER.error("consumer is start !!! groupName:{},topics:{},namesrvAddr:{}",groupName,topics,namesrvAddr,e);
    throw new RocketMQException(e);
    }
    return consumer;
    }
    }
    4、生产端配置
    package com.example.demo.rocketmq.producer;

import com.alibaba.rocketmq.client.exception.MQClientException;
import com.alibaba.rocketmq.client.producer.DefaultMQProducer;
import com.example.demo.rocketmq.constants.RocketMQErrorEnum;
import com.example.demo.rocketmq.exception.RocketMQException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.SpringBootConfiguration;
import org.springframework.context.annotation.Bean;
import org.springframework.util.StringUtils;

/**

  • 生产者配置

  • .

  • Copyright: Copyright (c) 2017 zteits

  • @ClassName: MQProducerConfiguration

  • @Description:

  • @version: v1.0.0

  • @author: zhaowg

  • @date: 2018年3月2日 下午11:44:36

  • Modification History:

  • Date Author Version Description
    ---------------------------------------------------------

  • 2018年3月2日 zhaowg v1.0.0 创建
    /
    @SpringBootConfiguration
    public class MQProducerConfiguration {
    public static final Logger LOGGER = LoggerFactory.getLogger(MQProducerConfiguration.class);
    /
    *

    • 发送同一类消息的设置为同一个group,保证唯一,默认不需要设置,rocketmq会使用ip@pid(pid代表jvm名字)作为唯一标示
      /
      @Value("{rocketmq.producer.groupName}") private String groupName; @Value("{rocketmq.producer.namesrvAddr}")
      private String namesrvAddr;
      /
      *
    • 消息最大大小,默认4M
      /
      @Value("${rocketmq.producer.maxMessageSize}")
      private Integer maxMessageSize ;
      /
      *
    • 消息发送超时时间,默认3秒
      /
      @Value("${rocketmq.producer.sendMsgTimeout}")
      private Integer sendMsgTimeout;
      /
      *
    • 消息发送失败重试次数,默认2次
      */
      @Value("${rocketmq.producer.retryTimesWhenSendFailed}")
      private Integer retryTimesWhenSendFailed;

    @Bean
    public DefaultMQProducer getRocketMQProducer() throws RocketMQException {
    if (StringUtils.isEmpty(this.groupName)) {
    throw new RocketMQException(RocketMQErrorEnum.PARAMM_NULL,"groupName is blank",false);
    }
    if (StringUtils.isEmpty(this.namesrvAddr)) {
    throw new RocketMQException(RocketMQErrorEnum.PARAMM_NULL,"nameServerAddr is blank",false);
    }
    DefaultMQProducer producer;
    producer = new DefaultMQProducer(this.groupName);
    producer.setNamesrvAddr(this.namesrvAddr);
    //如果需要同一个jvm中不同的producer往不同的mq集群发送消息,需要设置不同的instanceName
    //producer.setInstanceName(instanceName);
    if(this.maxMessageSize!=null){
    producer.setMaxMessageSize(this.maxMessageSize);
    }
    if(this.sendMsgTimeout!=null){
    producer.setSendMsgTimeout(this.sendMsgTimeout);
    }
    //如果发送消息失败,设置重试次数,默认为2次
    if(this.retryTimesWhenSendFailed!=null){
    producer.setRetryTimesWhenSendFailed(this.retryTimesWhenSendFailed);
    }

     try {
         producer.start();
         
         LOGGER.info(String.format("producer is start ! groupName:[%s],namesrvAddr:[%s]"
                 , this.groupName, this.namesrvAddr));
     } catch (MQClientException e) {
         LOGGER.error(String.format("producer is error {}"
                 , e.getMessage(),e));
         throw new RocketMQException(e);
     }
     return producer;
    

    }
    }

5.监听配置
package com.example.demo.rocketmq.consumer.processor;

import com.alibaba.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import com.alibaba.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import com.alibaba.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import com.alibaba.rocketmq.common.message.MessageExt;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
import org.springframework.util.CollectionUtils;

import java.util.List;

/**

  • 消费者消费消息路由

  • .

  • Copyright: Copyright (c) 2017 zteits

  • @ClassName: RocketMQMessageListenerConcurrentlyProcessor

  • @Description:

  • @version: v1.0.0

  • @author: zhaowg

  • @date: 2018年2月28日 上午11:12:32

  • Modification History:

  • Date Author Version Description
    ---------------------------------------------------------

  • 2018年2月28日 zhaowg v1.0.0 创建
    */
    @Component
    public class MQConsumeMsgListenerProcessor implements MessageListenerConcurrently{
    private static final Logger logger = LoggerFactory.getLogger(MQConsumeMsgListenerProcessor.class);

    /**

    • 默认msgs里只有一条消息,可以通过设置consumeMessageBatchMaxSize参数来批量接收消息
    • 不要抛异常,如果没有return CONSUME_SUCCESS ,consumer会重新消费该消息,直到return CONSUME_SUCCESS
      */
      @Override
      public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
      if(CollectionUtils.isEmpty(msgs)){
      logger.info("接受到的消息为空,不处理,直接返回成功");
      return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
      }
      MessageExt messageExt = msgs.get(0);
      logger.info("接受到的消息为:"+messageExt.toString());
      if(messageExt.getTopic().equals("你的Topic")){
      if(messageExt.getTags().equals("你的Tag")){
      //TODO 判断该消息是否重复消费(RocketMQ不保证消息不重复,如果你的业务需要保证严格的不重复消息,需要你自己在业务端去重)
      //TODO 获取该消息重试次数
      int reconsume = messageExt.getReconsumeTimes();
      if(reconsume ==3){//消息已经重试了3次,如果不需要再次消费,则返回成功
      return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
      }
      //TODO 处理对应的业务逻辑
      }
      }
      // 如果没有return success ,consumer会重新消费该消息,直到return success
      return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
      }
      }
      6、进行测试
      package com.example.demo;

import com.alibaba.rocketmq.client.exception.MQBrokerException;
import com.alibaba.rocketmq.client.exception.MQClientException;
import com.alibaba.rocketmq.client.producer.DefaultMQProducer;
import com.alibaba.rocketmq.client.producer.SendResult;
import com.alibaba.rocketmq.common.message.Message;
import com.alibaba.rocketmq.remoting.exception.RemotingException;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.junit4.SpringRunner;

@RunWith(SpringRunner.class)
@SpringBootTest
public class DemoApplicationTests {

private static final Logger logger = LoggerFactory.getLogger(DemoApplicationTests.class);

/**使用RocketMq的生产者*/
@Autowired
private DefaultMQProducer defaultMQProducer;

/**
 * 发送消息
 *
 * 2018年3月3日 zhaowg
 * @throws InterruptedException
 * @throws MQBrokerException
 * @throws RemotingException
 * @throws MQClientException
 */
@Test
public void send() throws MQClientException, RemotingException, MQBrokerException, InterruptedException{
    String msg = "demo msg test";
    logger.info("开始发送消息:"+msg);
    Message sendMsg = new Message("DemoTopic","DemoTag",msg.getBytes());
    //默认3秒超时
    SendResult sendResult = defaultMQProducer.send(sendMsg);
    logger.info("消息发送响应信息:"+sendResult.toString());
}

}

©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 204,445评论 6 478
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 85,889评论 2 381
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 151,047评论 0 337
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 54,760评论 1 276
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 63,745评论 5 367
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 48,638评论 1 281
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 38,011评论 3 398
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 36,669评论 0 258
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 40,923评论 1 299
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 35,655评论 2 321
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 37,740评论 1 330
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 33,406评论 4 320
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 38,995评论 3 307
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 29,961评论 0 19
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 31,197评论 1 260
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 45,023评论 2 350
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 42,483评论 2 342

推荐阅读更多精彩内容