SpringBoot集成RocketMQ

1. 導(dǎo)入依賴

compile 'org.apache.rocketmq:rocketmq-client:4.5.2'

image.png

2. 編寫application.yml配置

image.png

3. 引入配置信息

為了方便,在這里消費(fèi)者和生產(chǎn)者都放在一個(gè)項(xiàng)目里

引入生產(chǎn)者配置信息

package utry.config;

import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Configuration;

/**
 * @author
 * 消息生產(chǎn)者配置信息
 */
@Configuration
@ConfigurationProperties(prefix = "rocketmq.producer")
public class ProducerPropertiesConfig {

    @Value("${namesrvAddr}")
    private String namesrvAddr;

    private String groupName;

    private Integer maxMessageSize;

    private Integer sendMsgTimeout;

    private Integer retryTimesWhenSendFailed;

    public String getNamesrvAddr() {
        return namesrvAddr;
    }

    public void setNamesrvAddr(String namesrvAddr) {
        this.namesrvAddr = namesrvAddr;
    }

    public String getGroupName() {
        return groupName;
    }

    public void setGroupName(String groupName) {
        this.groupName = groupName;
    }

    public Integer getMaxMessageSize() {
        return maxMessageSize;
    }

    public void setMaxMessageSize(Integer maxMessageSize) {
        this.maxMessageSize = maxMessageSize;
    }

    public Integer getSendMsgTimeout() {
        return sendMsgTimeout;
    }

    public void setSendMsgTimeout(Integer sendMsgTimeout) {
        this.sendMsgTimeout = sendMsgTimeout;
    }

    public Integer getRetryTimesWhenSendFailed() {
        return retryTimesWhenSendFailed;
    }

    public void setRetryTimesWhenSendFailed(Integer retryTimesWhenSendFailed) {
        this.retryTimesWhenSendFailed = retryTimesWhenSendFailed;
    }

    @Override
    public String toString() {
        return "ProducerConfig [namesrvAddr=" + namesrvAddr + ", groupName=" + groupName + "]";
    }
}

編寫生產(chǎn)者



import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

/**
 * @author
 * 消息生產(chǎn)者
 */
@Configuration
public class ProducerConfigure {

    Logger logger = LoggerFactory.getLogger(ProducerConfigure.class);

    @Autowired
    private ProducerPropertiesConfig producerPropertiesConfig;

    @Bean
    public DefaultMQProducer defaultProducer() throws MQClientException {
        logger.info(producerPropertiesConfig.toString());
        logger.info("defaultProducer 正在創(chuàng)建---------------------------------------");
        DefaultMQProducer producer = new DefaultMQProducer(producerPropertiesConfig.getGroupName());
        producer.setNamesrvAddr(producerPropertiesConfig.getNamesrvAddr());
        producer.setVipChannelEnabled(false);
        //其他屬性自行設(shè)置矾瑰,這里才用默認(rèn)
        producer.start();
        logger.info("rocketmq producer server開啟成功---------------------------------.");
        return producer;
    }
}

引入消費(fèi)者配置信息

package utry.config;

import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Configuration;

/**
 * @author
 * 消費(fèi)者屬性配置類
 */
@ConfigurationProperties(prefix = "rocketmq.consumer")
@Configuration
public class ConsumerPropertiesConfig {
    private String groupName;

    @Value("${namesrvAddr}")
    private String namesrvAddr;

    public String getGroupName() {
        return groupName;
    }

    public void setGroupName(String groupName) {
        this.groupName = groupName;
    }

    public String getNamesrvAddr() {
        return namesrvAddr;
    }

    public void setNamesrvAddr(String namesrvAddr) {
        this.namesrvAddr = namesrvAddr;
    }

    @Override
    public String toString() {
        return "ConsumerConfig [groupName=" + groupName + ", namesrvAddr=" + namesrvAddr + "]";
    }

}

編寫消費(fèi)者
先編寫一個(gè)抽象類,再寫具體的實(shí)現(xiàn)

  1. 編寫抽象類
package utry.config;

import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
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.message.MessageExt;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;

import java.util.List;

/**
 * @author
 * 抽象的消息消費(fèi)者
 */
@Service
public abstract class DefaultConsumerConfigure {
    Logger log = LoggerFactory.getLogger(DefaultConsumerConfigure.class);

    @Autowired
    private ConsumerPropertiesConfig consumerConfig;

    /**
     * 開啟消費(fèi)者監(jiān)聽服務(wù)
     * @param topic
     * @param tag
     * @throws MQClientException
     */
    public void listener(String topic, String tag) throws MQClientException {
        log.info("開啟" + topic + ":" + tag + "消費(fèi)者-------------------");
        log.info(consumerConfig.toString());

        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(consumerConfig.getGroupName());

        consumer.setNamesrvAddr(consumerConfig.getNamesrvAddr());

        consumer.subscribe(topic, tag);

        // 開啟內(nèi)部類實(shí)現(xiàn)監(jiān)聽
        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> messageExtList, ConsumeConcurrentlyContext context) {
                return DefaultConsumerConfigure.this.dealBody(messageExtList);
            }
        });
        consumer.start();

        log.info("rocketmq啟動(dòng)成功---------------------------------------");

    }

    /**
     * 處理body的業(yè)務(wù)
     * @param messageExtList
     * @return
     */
    public abstract ConsumeConcurrentlyStatus dealBody(List<MessageExt> messageExtList);

}

  1. 編寫消費(fèi)者
package utry.config;

import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.common.message.MessageExt;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.ApplicationListener;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.event.ContextRefreshedEvent;

import java.util.List;

/**
 * @author
 * 消息消費(fèi)者
 */
@Configuration
public class CustomConsumer extends DefaultConsumerConfigure implements ApplicationListener<ContextRefreshedEvent> {

    Logger log = LoggerFactory.getLogger(CustomConsumer.class);

    @Override
    public void onApplicationEvent(ContextRefreshedEvent contextRefreshedEvent) {
        try {
            super.listener("t_TopicTest", "Tag1");
        } catch (MQClientException e) {
            log.error("消費(fèi)者監(jiān)聽器啟動(dòng)失敗", e);
        }

    }

    @Override
    public ConsumeConcurrentlyStatus dealBody(List<MessageExt> messageExtList) {
        log.info("接收到消息");
        for (MessageExt msg : messageExtList) {
            try {
                String msgStr = new String(msg.getBody(), "utf-8");
                log.info(msgStr);
            } catch (Exception e) {
                log.error("body轉(zhuǎn)字符串解析失敗");
            }
        }
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    }

}

編寫Controller測(cè)試

package utry.controller;

import com.alibaba.fastjson.JSON;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import utry.config.CustomConsumer;

@RestController
public class ProducerController {

    Logger log = LoggerFactory.getLogger(CustomConsumer.class);

    @Autowired
    private DefaultMQProducer producer;


    @GetMapping("/msg/product")
    public void test(String info) throws Exception {
        Message message = new Message("t_TopicTest", "Tag1", "12345", info.getBytes());
        // 這里用到了這個(gè)mq的異步處理,類似ajax,可以得到發(fā)送到mq的情況壶愤,并做相應(yīng)的處理
        // 不過要注意的是這個(gè)是異步的
        producer.send(message, new SendCallback() {
            @Override
            public void onSuccess(SendResult sendResult) {
                log.info("傳輸成功");
                log.info(JSON.toJSONString(sendResult));
            }

            @Override
            public void onException(Throwable e) {
                log.error("傳輸失敗", e);
            }
        });
    }

}

github:

?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請(qǐng)聯(lián)系作者
  • 序言:七十年代末揭绑,一起剝皮案震驚了整個(gè)濱河市,隨后出現(xiàn)的幾起案子库说,更是在濱河造成了極大的恐慌,老刑警劉巖片择,帶你破解...
    沈念sama閱讀 212,816評(píng)論 6 492
  • 序言:濱河連續(xù)發(fā)生了三起死亡事件潜的,死亡現(xiàn)場離奇詭異,居然都是意外死亡字管,警方通過查閱死者的電腦和手機(jī)啰挪,發(fā)現(xiàn)死者居然都...
    沈念sama閱讀 90,729評(píng)論 3 385
  • 文/潘曉璐 我一進(jìn)店門,熙熙樓的掌柜王于貴愁眉苦臉地迎上來嘲叔,“玉大人亡呵,你說我怎么就攤上這事×蚋辏” “怎么了锰什?”我有些...
    開封第一講書人閱讀 158,300評(píng)論 0 348
  • 文/不壞的土叔 我叫張陵,是天一觀的道長。 經(jīng)常有香客問我歇由,道長卵牍,這世上最難降的妖魔是什么? 我笑而不...
    開封第一講書人閱讀 56,780評(píng)論 1 285
  • 正文 為了忘掉前任沦泌,我火速辦了婚禮糊昙,結(jié)果婚禮上,老公的妹妹穿的比我還像新娘谢谦。我一直安慰自己释牺,他們只是感情好,可當(dāng)我...
    茶點(diǎn)故事閱讀 65,890評(píng)論 6 385
  • 文/花漫 我一把揭開白布回挽。 她就那樣靜靜地躺著没咙,像睡著了一般。 火紅的嫁衣襯著肌膚如雪千劈。 梳的紋絲不亂的頭發(fā)上祭刚,一...
    開封第一講書人閱讀 50,084評(píng)論 1 291
  • 那天,我揣著相機(jī)與錄音墙牌,去河邊找鬼涡驮。 笑死,一個(gè)胖子當(dāng)著我的面吹牛喜滨,可吹牛的內(nèi)容都是我干的捉捅。 我是一名探鬼主播,決...
    沈念sama閱讀 39,151評(píng)論 3 410
  • 文/蒼蘭香墨 我猛地睜開眼虽风,長吁一口氣:“原來是場噩夢(mèng)啊……” “哼棒口!你這毒婦竟也來了?” 一聲冷哼從身側(cè)響起辜膝,我...
    開封第一講書人閱讀 37,912評(píng)論 0 268
  • 序言:老撾萬榮一對(duì)情侶失蹤无牵,失蹤者是張志新(化名)和其女友劉穎,沒想到半個(gè)月后厂抖,有當(dāng)?shù)厝嗽跇淞掷锇l(fā)現(xiàn)了一具尸體合敦,經(jīng)...
    沈念sama閱讀 44,355評(píng)論 1 303
  • 正文 獨(dú)居荒郊野嶺守林人離奇死亡,尸身上長有42處帶血的膿包…… 初始之章·張勛 以下內(nèi)容為張勛視角 年9月15日...
    茶點(diǎn)故事閱讀 36,666評(píng)論 2 327
  • 正文 我和宋清朗相戀三年验游,在試婚紗的時(shí)候發(fā)現(xiàn)自己被綠了。 大學(xué)時(shí)的朋友給我發(fā)了我未婚夫和他白月光在一起吃飯的照片保檐。...
    茶點(diǎn)故事閱讀 38,809評(píng)論 1 341
  • 序言:一個(gè)原本活蹦亂跳的男人離奇死亡耕蝉,死狀恐怖,靈堂內(nèi)的尸體忽然破棺而出夜只,到底是詐尸還是另有隱情垒在,我是刑警寧澤,帶...
    沈念sama閱讀 34,504評(píng)論 4 334
  • 正文 年R本政府宣布扔亥,位于F島的核電站场躯,受9級(jí)特大地震影響谈为,放射性物質(zhì)發(fā)生泄漏。R本人自食惡果不足惜踢关,卻給世界環(huán)境...
    茶點(diǎn)故事閱讀 40,150評(píng)論 3 317
  • 文/蒙蒙 一伞鲫、第九天 我趴在偏房一處隱蔽的房頂上張望。 院中可真熱鬧签舞,春花似錦秕脓、人聲如沸。這莊子的主人今日做“春日...
    開封第一講書人閱讀 30,882評(píng)論 0 21
  • 文/蒼蘭香墨 我抬頭看了看天上的太陽。三九已至搂鲫,卻和暖如春傍药,著一層夾襖步出監(jiān)牢的瞬間,已是汗流浹背魂仍。 一陣腳步聲響...
    開封第一講書人閱讀 32,121評(píng)論 1 267
  • 我被黑心中介騙來泰國打工拐辽, 沒想到剛下飛機(jī)就差點(diǎn)兒被人妖公主榨干…… 1. 我叫王不留,地道東北人蓄诽。 一個(gè)月前我還...
    沈念sama閱讀 46,628評(píng)論 2 362
  • 正文 我出身青樓薛训,卻偏偏與公主長得像,于是被迫代替她去往敵國和親仑氛。 傳聞我的和親對(duì)象是個(gè)殘疾皇子乙埃,可洞房花燭夜當(dāng)晚...
    茶點(diǎn)故事閱讀 43,724評(píng)論 2 351