springBoot整合RocketMQ及坑的示例代碼
版本:
- JDK:1.8
- springBoot:1.5.10
- rocketMQ:4.2.0
pom 配置:
<parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>1.5.10.RELEASE</version> </parent> <dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-client</artifactId> <version>4.2.0</version> </dependency>
application.properties 配置:
# 消費(fèi)者的組名 apache.rocketmq.consumer.PushConsumer=PushConsumer # 生產(chǎn)者的組名 apache.rocketmq.producer.producerGroup=Producer # NameServer地址 apache.rocketmq.namesrvAddr=localhost:9876
java代碼:
生產(chǎn)者
package test.config.rocketmq;
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.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import org.springframework.util.StopWatch;
import javax.annotation.PostConstruct;
@Component
public class RocketMQClient {
/**
* 生產(chǎn)者的組名
*/
@Value("${apache.rocketmq.producer.producerGroup}")
private String producerGroup;
/**
* NameServer 地址
*/
@Value("${apache.rocketmq.namesrvAddr}")
private String namesrvAddr;
@PostConstruct
public void defaultMQProducer() {
//生產(chǎn)者的組名
DefaultMQProducer producer = new DefaultMQProducer(producerGroup);
//指定NameServer地址,多個(gè)地址以 ; 隔開
producer.setNamesrvAddr(namesrvAddr);
producer.setVipChannelEnabled(false);
try {
/**
* Producer對(duì)象在使用之前必須要調(diào)用start初始化,初始化一次即可
* 注意:切記不可以在每次發(fā)送消息時(shí),都調(diào)用start方法
*/
producer.start();
//創(chuàng)建一個(gè)消息實(shí)例,包含 topic、tag 和 消息體
//如下:topic 為 "TopicTest",tag 為 "push"
Message message = new Message("TopicTest", "push", "發(fā)送消息----zhisheng-----".getBytes(RemotingHelper.DEFAULT_CHARSET));
StopWatch stop = new StopWatch();
stop.start();
for (int i = 0; i < 1; i++) {
SendResult result = producer.send(message);
System.out.println("發(fā)送響應(yīng):MsgId:" + result.getMsgId() + ",發(fā)送狀態(tài):" + result.getSendStatus());
}
stop.stop();
System.out.println("----------------發(fā)送一萬(wàn)條消息耗時(shí):" + stop.getTotalTimeMillis());
} catch (Exception e) {
e.printStackTrace();
} finally {
producer.shutdown();
}
}
}
消費(fèi)者:
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.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import javax.annotation.PostConstruct;
@Component
public class RocketMQServer {
/**
* 消費(fèi)者的組名
*/
@Value("${apache.rocketmq.consumer.PushConsumer}")
private String consumerGroup;
/**
* NameServer 地址
*/
@Value("${apache.rocketmq.namesrvAddr}")
private String namesrvAddr;
@PostConstruct
public void defaultMQPushConsumer() {
//消費(fèi)者的組名
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(consumerGroup);
//指定NameServer地址,多個(gè)地址以 ; 隔開
consumer.setNamesrvAddr(namesrvAddr);
consumer.setVipChannelEnabled(false);
try {
//訂閱PushTopic下Tag為push的消息
consumer.subscribe("TopicTest", "push");
//設(shè)置Consumer第一次啟動(dòng)是從隊(duì)列頭部開始消費(fèi)還是隊(duì)列尾部開始消費(fèi)
//如果非第一次啟動(dòng),那么按照上次消費(fèi)的位置繼續(xù)消費(fèi)
consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
consumer.registerMessageListener((MessageListenerConcurrently) (list, context) -> {
try {
for (MessageExt messageExt : list) {
System.out.println("messageExt: " + messageExt);//輸出消息內(nèi)容
String messageBody = new String(messageExt.getBody(), RemotingHelper.DEFAULT_CHARSET);
System.out.println("消費(fèi)響應(yīng):msgId : " + messageExt.getMsgId() + ", msgBody : " + messageBody);//輸出消息內(nèi)容
}
} catch (Exception e) {
e.printStackTrace();
return ConsumeConcurrentlyStatus.RECONSUME_LATER; //稍后再試
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; //消費(fèi)成功
});
consumer.start();
} catch (Exception e) {
e.printStackTrace();
}
}
}
掉坑總結(jié):
1.rocketMQ啟動(dòng)時(shí),命令不是 mqbroker -n 127.0.0.1:9876
正確應(yīng)該是:mqbroker -n 127.0.0.1:9876 butiautoCreateTopicEnable=true
否則會(huì)拋出:No route info of this topic, TopicTest
2.客戶端連接時(shí)拋出異常
org.apache.rocketmq.client.exception.MQClientException:
Send [3] times, still failed, cost [3180]ms, Topic: TopicTest, BrokersSent: \
[WIN-93CGO0S5G25, WIN-93CGO0S5G25, WIN-93CGO0S5G25]
解決方式兩種
1.producer.setVipChannelEnabled(false); 生產(chǎn)者和消費(fèi)者添加這行代買。
2.降rocketmq版本,降成3.2.6
關(guān)于spring.rocketmq.name-server的坑
看下圖:

注意:
如果你是SpringBoot2.0+的框架,或者是JDK10。
你需要將你自己的項(xiàng)目配置文件中的,spring.rocketmq.name-server改成
spring.rocketmq.nameServer。注意是nameServer。
不然就會(huì)報(bào)各種稀奇古怪的bug。
關(guān)于啟動(dòng)報(bào)內(nèi)存不足的錯(cuò)
在安裝啟動(dòng)Name Server和Broker的時(shí)候,一定要修改配置文件,不然內(nèi)存會(huì)爆炸。
Native memory allocation (mmap) failed to map 8589934592 bytes for committing reserved memory

將下面的配置文件根據(jù)你的需要改
我這里以前默認(rèn)是Xms4g,都是g,我修改到m就行了。
JAVA_OPT="${JAVA_OPT} -server -Xms256m -Xmx256m -Xmn128m -XX:MetaspaceSize=128m -XX:MaxMetaspaceSize=320m"
以上就是本文的全部?jī)?nèi)容,希望對(duì)大家的學(xué)習(xí)有所幫助,也希望大家多多支持腳本之家。
- 淺談Springboot整合RocketMQ使用心得
- Springboot RocketMq實(shí)現(xiàn)過程詳解
- Springboot詳解RocketMQ實(shí)現(xiàn)廣播消息流程
- 解決SpringBoot整合RocketMQ遇到的坑
- SpringBoot整合RocketMQ實(shí)現(xiàn)發(fā)送同步消息
- 解決springboot集成rocketmq關(guān)于tag的坑
- SpringBoot集成RocketMQ的使用示例
- 在SpringBoot中利用RocketMQ實(shí)現(xiàn)批量消息消費(fèi)功能
- SpringBoot項(xiàng)目嵌入RocketMQ的實(shí)現(xiàn)示例
- RocketMQ在Spring Boot上的基礎(chǔ)使用
相關(guān)文章
java web實(shí)現(xiàn)用戶權(quán)限管理
這篇文章主要介紹了java web實(shí)現(xiàn)用戶權(quán)限管理,設(shè)計(jì)并實(shí)現(xiàn)一套簡(jiǎn)單的權(quán)限管理功能,感興趣的小伙伴們可以參考一下2015-11-11
pdf2swf+flexpapers實(shí)現(xiàn)類似百度文庫(kù)pdf在線閱讀
這篇文章主要介紹了pdf2swf+flexpapers實(shí)現(xiàn)類似百度文庫(kù)pdf在線閱讀的相關(guān)資料,需要的朋友可以參考下2014-10-10

