国产成人精品久久免费动漫-国产成人精品天堂-国产成人精品区在线观看-国产成人精品日本-a级毛片无码免费真人-a级毛片毛片免费观看久潮喷

您的位置:首頁技術文章
文章詳情頁

Springboot RocketMq實現過程詳解

瀏覽:123日期:2023-05-16 13:09:35

首先,在虛擬機上安裝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 **/@Componentpublic 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 **/@Componentpublic 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 RocketMq實現過程詳解

以上就是本文的全部內容,希望對大家的學習有所幫助,也希望大家多多支持好吧啦網。

標簽: Spring
相關文章:
主站蜘蛛池模板: 91精品一区二区综合在线 | 男女性生活网站 | 一区二区三区中文国产亚洲 | 中文字幕综合在线 | 欧美成人性动漫在线观看 | 一区国严二区亚洲三区 | 国产精品美女久久久久网站 | 国产无套视频在线观看香蕉 | 亚洲一区二区视频 | 久草视频2 | 国产精品一二三区 | 台湾三级香港三级在线中文 | 国产久视频 | 欧美成人免费午夜全 | 四色6677最新永久网站 | 国产精品99精品久久免费 | 国产成人亚洲精品影院 | 老色99久久九九精品尤物 | 国产专区在线 | 蜜桃日本一道无卡不码高清 | 黄色毛片视频校园交易 | 国产日产欧美a级毛片 | 成人18视频在线观看 | 成人69视频在线观看免费 | 玖草| 直接在线观看的三级网址 | 国产精品免费视频能看 | 在线91精品国产免费 | 一级做a爰片久久毛片人呢 一级做a爰片久久毛片唾 | 日本不卡一区在线 | 香蕉久久久 | 韩国成人毛片aaa黄 韩国福利一区 | 欧美一区中文字幕 | 成人亚洲国产精品久久 | 中文字幕有码在线 | 一本色道久久88综合亚洲精品高清 | 国产成人在线视频观看 | 波多野结衣一区在线观看 | 亚洲精品一区二区不卡 | 国产日韩欧美精品在线 | 中文国产成人精品久久水 |