參考文章:https://blog.csdn.net/qq_18603599/article/details/81172866
生産者
/**
* @Function 模拟使用者消息發送
*/
@Component
public class UserProducer {
/**
* 生産者的組名
*/
@Value("${suning.rocketmq.producerGroup}")
private String producerGroup;
/**
* NameServer 位址
*/
@Value("${suning.rocketmq.namesrvaddr}")
private String namesrvAddr;
// @PostConstruct:它的作用是相當于servlet的Init功能,在對象的構造器執行完了,就會立馬調用該方法
@PostConstruct
public void produder() {
DefaultMQProducer producer = new DefaultMQProducer(producerGroup);
producer.setNamesrvAddr(namesrvAddr);
try {
producer.start();
for (int i = 0; i < 100; i++) {
UserContent userContent = new UserContent(String.valueOf(i),"abc"+i);
String jsonstr = JSON.toJSONString(userContent);
System.out.println("發送消息:"+jsonstr);
Message message = new Message("user-topic", "user-tag", jsonstr.getBytes(RemotingHelper.DEFAULT_CHARSET));
SendResult result = producer.send(message);
System.err.println("發送響應:MsgId:" + result.getMsgId() + ",發送狀态:" + result.getSendStatus());
}
} catch (Exception e) {
e.printStackTrace();
} finally {
producer.shutdown();
}
}
}
消費者:
@Component
public class UserConsumer {
/**
* 消費者的組名
*/
@Value("${suning.rocketmq.conumerGroup}")
private String consumerGroup;
/**
* NameServer 位址
*/
@Value("${suning.rocketmq.namesrvaddr}")
private String namesrvAddr;
@PostConstruct
public void consumer() {
System.err.println("init defaultMQPushConsumer");
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(consumerGroup);
consumer.setNamesrvAddr(namesrvAddr);
try {
consumer.subscribe("user-topic", "user-tag");
consumer.registerMessageListener((MessageListenerConcurrently) (list, context) -> {
try {
for (MessageExt messageExt : list) {
System.err.println("消費消息: " + new String(messageExt.getBody()));//輸出消息内容
}
} catch (Exception e) {
e.printStackTrace();
return ConsumeConcurrentlyStatus.RECONSUME_LATER; //稍後再試
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; //消費成功
});
consumer.start();
} catch (Exception e) {
e.printStackTrace();
}
}
}
配置檔案
# 生産者的組名
suning.rocketmq.producerGroup=user-group
# 消費者的組名
suning.rocketmq.conumerGroup=user-group
# NameServer位址
suning.rocketmq.namesrvaddr=127.0.0.1:9876