Springboot RocketMq實(shí)現(xiàn)過(guò)程詳解
首先,在虛擬機(jī)上安裝rocketmq和rocketMq可視化控制,安裝不做描述。
1、pom.xml文件添加依賴(lài)
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)者類(lèi)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 **/@Componentpublic 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)者類(lè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 **/@Componentpublic 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)聽(tīng)信息,渠道信息 */ 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)聽(tīng)中消費(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中編寫(xiě)發(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'}
通過(guò)rocketMq可視化控制查看:
以上就是本文的全部?jī)?nèi)容,希望對(duì)大家的學(xué)習(xí)有所幫助,也希望大家多多支持好吧啦網(wǎng)。
相關(guān)文章:
1. XML在語(yǔ)音合成中的應(yīng)用2. HTTP協(xié)議常用的請(qǐng)求頭和響應(yīng)頭響應(yīng)詳解說(shuō)明(學(xué)習(xí))3. 不要在HTML中濫用div4. ASP將數(shù)字轉(zhuǎn)中文數(shù)字(大寫(xiě)金額)的函數(shù)5. .NET Framework各版本(.NET2.0 3.0 3.5 4.0)區(qū)別6. jscript與vbscript 操作XML元素屬性的代碼7. HTML5實(shí)戰(zhàn)與剖析之觸摸事件(touchstart、touchmove和touchend)8. php使用正則驗(yàn)證密碼字段的復(fù)雜強(qiáng)度原理詳細(xì)講解 原創(chuàng)9. ASP基礎(chǔ)入門(mén)第四篇(腳本變量、函數(shù)、過(guò)程和條件語(yǔ)句)10. XML入門(mén)的常見(jiàn)問(wèn)題(三)
