RocketMQ的兩種消費(fèi)模式詳解
1、添加依賴
<dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-client</artifactId> <version>4.4.0</version> </dependency> <dependency> <groupId>com.alibaba</groupId> <artifactId>fastjson</artifactId> <version>1.2.3</version> </dependency> <dependencies> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> </dependency> </dependencies>
2、消費(fèi)模式
RocketMQ主要提供了兩種消費(fèi)模式:集群消費(fèi)以及廣播消費(fèi)。我們只需要在定義消費(fèi)者的時(shí)候通過setMessageModel(MessageModel.XXX)
// 設(shè)置消費(fèi)模型,集群還是廣播,默認(rèn)為集群 CLUSTERING-集群,BROADCASTING-廣播 mqPushConsumer.setMessageModel(MessageModel.CLUSTERING); mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);
方法就可以指定是集群還是廣播式消費(fèi),默認(rèn)是集群消費(fèi)模式,即每個(gè)Consumer Group中的Consumer均攤所有的消息。
3、集群消費(fèi)
3.1 生產(chǎn)者
package com.shucha.deveiface.biz.mq.producer;
import com.alibaba.fastjson.JSON;
import com.shucha.deveiface.biz.model.User;
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 java.util.Date;
/**
* @author tqf
* @Description 生產(chǎn)者
* @Version 1.0
* @since 2022-07-12 14:50
*/
public class MQProducer {
public static void main(String[] args) throws MQClientException{
producerSendMessage();
}
/**
* 生產(chǎn)消息方法
* @throws MQClientException
*/
public static void producerSendMessage() throws MQClientException {
// 創(chuàng)建DefaultMQProducer類并設(shè)定生產(chǎn)者名稱
DefaultMQProducer mqProducer = new DefaultMQProducer("producer-group-test");
// 設(shè)置NameServer地址,如果是集群的話,使用分號;分隔開
mqProducer.setNamesrvAddr("127.0.0.1:9876");
// 消息最大長度 默認(rèn)4M
mqProducer.setMaxMessageSize(4096);
// 發(fā)送消息超時(shí)時(shí)間,默認(rèn)3000
mqProducer.setSendMsgTimeout(3000);
// 發(fā)送消息失敗重試次數(shù),默認(rèn)2
mqProducer.setRetryTimesWhenSendAsyncFailed(2);
// 啟動(dòng)消息生產(chǎn)者
mqProducer.start();
try {
// 循環(huán)十次,發(fā)送十條消息
for (int i = 1; i <= 10; i++) {
User user = new User();
user.setId((long)i);
user.setAge(i);
user.setUserName("姓名"+i);
user.setCreateTime(new Date());
// String msg = "這是第" + i + "條消息測試";
String msg = JSON.toJSONString(user);
// 創(chuàng)建消息,并指定Topic(主題),Tag(標(biāo)簽)和消息內(nèi)容
Message message = new Message("TOPIC_TEST", "", msg.getBytes(RemotingHelper.DEFAULT_CHARSET));
// 發(fā)送同步消息到一個(gè)Broker,可以通過sendResult返回消息是否成功送達(dá)
SendResult sendResult = mqProducer.send(message);
// mqProducer.sendOneway(message);
// 消息id
/*System.out.println(sendResult.getMsgId());
// 隊(duì)列信息
System.out.println(sendResult.getMessageQueue());
// 發(fā)送結(jié)果
System.out.println(sendResult.getSendStatus());
// 下一個(gè)要消費(fèi)的消息的偏移量
System.out.println(sendResult.getOffsetMsgId());
// 隊(duì)列消息偏移量
System.out.println(sendResult.getQueueOffset());*/
System.out.println(sendResult);
}
} catch (Exception e) {
e.printStackTrace();
System.out.println("生產(chǎn)消息異常!");
}
// 如果不再發(fā)送消息,關(guān)閉Producer實(shí)例
mqProducer.shutdown();
}
}User用戶測試實(shí)體類
package com.shucha.deveiface.biz.model;
import com.fasterxml.jackson.annotation.JsonFormat;
import com.sdy.common.utils.DateUtil;
import io.swagger.annotations.ApiModelProperty;
import lombok.Data;
import java.util.Date;
/**
* @author tqf
* @Description
* @Version 1.0
* @since 2022-04-07 13:53
*/
@Data
public class User {
/**
* 主鍵ID
*/
private Long id;
/**
*用戶名
*/
private String userName;
/**
* 用戶密碼
*/
private String passWord;
/**
* 年齡
*/
private Integer age;
/**
* 性別(0-男,1-女,2-未知)
*/
private Integer sex;
/**
* 創(chuàng)建時(shí)間
*/
@JsonFormat(pattern = DateUtil.DATETIME_FORMAT)
private Date createTime;
}3.2 消費(fèi)者A
package com.shucha.deveiface.biz.mq.consumer;
import com.alibaba.fastjson.JSON;
import com.shucha.deveiface.biz.constants.MqConstants;
import com.shucha.deveiface.biz.model.User;
import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
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.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import java.nio.charset.Charset;
import java.util.List;
/**
* @author tqf
* @Description 消費(fèi)者A
* @Version 1.0
* @since 2022-07-12 14:37
*/
public class ConsumerA {
public static void main(String[] args) throws MQClientException {
ConsumerA();
}
public static void ConsumerA() throws MQClientException {
// 創(chuàng)建DefaultMQPushConsumer類并設(shè)定消費(fèi)者名稱
DefaultMQPushConsumer mqPushConsumer = new DefaultMQPushConsumer(MqConstants.ConsumerGroup.CONSUMER_GROUP1);
// DefaultMQPullConsumer pullConsumer = new DefaultMQPullConsumer(MqConstants.ConsumerGroup.CONSUMER_GROUP1);
// 設(shè)置NameServer地址,如果是集群的話,使用分號;分隔開
mqPushConsumer.setNamesrvAddr("127.0.0.1:9876");
// pullConsumer.setNamesrvAddr("127.0.0.1:9876");
// 設(shè)置Consumer第一次啟動(dòng)是從隊(duì)列頭部開始消費(fèi)還是隊(duì)列尾部開始消費(fèi)
// 如果不是第一次啟動(dòng),那么按照上次消費(fèi)的位置繼續(xù)消費(fèi)
mqPushConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
// 設(shè)置消費(fèi)模型,集群還是廣播,默認(rèn)為集群 CLUSTERING-集群 BROADCASTING-廣播
mqPushConsumer.setMessageModel(MessageModel.CLUSTERING);
// mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);
// 消費(fèi)者最小線程量
mqPushConsumer.setConsumeThreadMin(5);
// 消費(fèi)者最大線程量
mqPushConsumer.setConsumeThreadMax(10);
// 設(shè)置一次消費(fèi)消息的條數(shù),默認(rèn)是1
mqPushConsumer.setConsumeMessageBatchMaxSize(1);
// 訂閱一個(gè)或者多個(gè)Topic,以及Tag來過濾需要消費(fèi)的消息,如果訂閱該主題下的所有tag,則使用*
mqPushConsumer.subscribe("TOPIC_TEST", "*");
// 注冊回調(diào)實(shí)現(xiàn)類來處理從broker拉取回來的消息
mqPushConsumer.registerMessageListener(new MessageListenerConcurrently() {
// 監(jiān)聽類實(shí)現(xiàn)MessageListenerConcurrently接口即可,重寫consumeMessage方法接收數(shù)據(jù)
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgList, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
for (MessageExt msg : msgList) {
String msgBody = new String(msg.getBody(), Charset.forName(RemotingHelper.DEFAULT_CHARSET));
System.out.println("消費(fèi)者A接收到消息:" + "===== " + msgBody);
/*User user = JSON.parseObject(msgBody, User.class);
System.out.println("消費(fèi)者A接收到消息:" + "===== " + user.getId());*/
}
/*MessageExt messageExt = msgList.get(0);
String body = new String(messageExt.getBody(), StandardCharsets.UTF_8);
System.out.println("消費(fèi)者A接收到消息: " + messageExt.toString() + "---消息內(nèi)容為:" + body);*/
// 標(biāo)記該消息已經(jīng)被成功消費(fèi)
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
// 啟動(dòng)消費(fèi)者實(shí)例
mqPushConsumer.start();
System.out.println("ConsumerA Started.");
}
}
3.3 消費(fèi)者B
package com.shucha.deveiface.biz.mq.consumer;
import com.shucha.deveiface.biz.constants.MqConstants;
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.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
import java.util.List;
/**
* @author tqf
* @Description 消費(fèi)者B
* @Version 1.0
* @since 2022-07-12 14:40
*/
public class ConsumerB {
public static void main(String[] args) throws MQClientException {
ConsumerB();
}
public static void ConsumerB() throws MQClientException {
// 創(chuàng)建DefaultMQPushConsumer類并設(shè)定消費(fèi)者名稱
DefaultMQPushConsumer mqPushConsumer = new DefaultMQPushConsumer(MqConstants.ConsumerGroup.CONSUMER_GROUP1);
// 設(shè)置NameServer地址,如果是集群的話,使用分號;分隔開
mqPushConsumer.setNamesrvAddr("127.0.0.1:9876");
// 設(shè)置Consumer第一次啟動(dòng)是從隊(duì)列頭部開始消費(fèi)還是隊(duì)列尾部開始消費(fèi)
// 如果不是第一次啟動(dòng),那么按照上次消費(fèi)的位置繼續(xù)消費(fèi)
mqPushConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
// 設(shè)置消費(fèi)模型,集群還是廣播,默認(rèn)為集群 CLUSTERING-集群 BROADCASTING-廣播
// mqPushConsumer.setMessageModel(MessageModel.CLUSTERING);
mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);
// 消費(fèi)者最小線程量
mqPushConsumer.setConsumeThreadMin(5);
// 消費(fèi)者最大線程量
mqPushConsumer.setConsumeThreadMax(10);
// 設(shè)置一次消費(fèi)消息的條數(shù),默認(rèn)是1
mqPushConsumer.setConsumeMessageBatchMaxSize(1);
// 訂閱一個(gè)或者多個(gè)Topic,以及Tag來過濾需要消費(fèi)的消息,如果訂閱該主題下的所有tag,則使用*
mqPushConsumer.subscribe("TOPIC_TEST", "*");
// 注冊回調(diào)實(shí)現(xiàn)類來處理從broker拉取回來的消息
mqPushConsumer.registerMessageListener(new MessageListenerConcurrently() {
// 監(jiān)聽類實(shí)現(xiàn)MessageListenerConcurrently接口即可,重寫consumeMessage方法接收數(shù)據(jù)
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgList, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
/* MessageExt messageExt = msgList.get(0);
String body = new String(messageExt.getBody(), StandardCharsets.UTF_8);
System.out.println("消費(fèi)者B接收到消息: " + messageExt.toString() + "---消息內(nèi)容為:" + body);*/
for (MessageExt msg : msgList) {
System.out.println("消費(fèi)者B接收到消息:" + "===== " + new String(msg.getBody()));
// String msgBody = new String(msg.getBody(), Charset.forName(RemotingHelper.DEFAULT_CHARSET));
}
// 標(biāo)記該消息已經(jīng)被成功消費(fèi)
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
// 啟動(dòng)消費(fèi)者實(shí)例
mqPushConsumer.start();
System.out.println("ConsumerB Started.");
}
}
可以看到, 生產(chǎn)者發(fā)送了10條消息,ConsumerA與ConsumerB屬于同一個(gè)消費(fèi)者組,集群消費(fèi)模式下每個(gè)消費(fèi)者攤分消費(fèi)所有消息。注意,兩個(gè)消費(fèi)者的ConsumerGroup組名需要一致,才算是同一個(gè)消費(fèi)者組。
簡單總結(jié)一下:
1、在Rocket集群消費(fèi)模式下,(訂閱)同一個(gè)主題(Topic)下的消息,對于不同的消費(fèi)者組是一種“廣播形式”,即每個(gè)消費(fèi)者組的都會消費(fèi)消息。
2、在Rocket集群消費(fèi)模式下,(訂閱)同一個(gè)主題(Topic)下的消息,對于相同的消費(fèi)者組的消費(fèi)者而言是一種集群模式,即同一個(gè)消費(fèi)者組內(nèi)的所有消費(fèi)者均分消息并消費(fèi)。
4、廣播消費(fèi)
一條消息被多個(gè) Consumer 消費(fèi),即使這些 Consumer 屬于同一個(gè) Consumer Group,消息也會被 Consumer Group 中的每個(gè) Consumer 都消費(fèi)一次,廣播消費(fèi)中的 Consumer Group 概念可以認(rèn)為在消息劃分方面無意 義。
使用方法:setMessageModel(MessageModel.BROADCASTING)
將前面的消費(fèi)者A和消費(fèi)者B的集群模式代碼設(shè)置為如下
mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);
重新啟動(dòng)生成者和2個(gè)消費(fèi)者
4.1 生產(chǎn)者數(shù)據(jù)

4.2 消費(fèi)者A

4.3 消費(fèi)者B

可以看到生產(chǎn)者發(fā)送了10條消息,ConsumerA與ConsumerB屬于同一個(gè)消費(fèi)者組,廣播模式下每個(gè)消費(fèi)者都會全量消費(fèi)所有消息 。
- 集群消費(fèi):任何一條消息只需要被消費(fèi)者集群中任意一個(gè)消費(fèi)者處理。
- 廣播消費(fèi):每條消息被推送給消費(fèi)者集群中的所有注冊消費(fèi)者,保證消息被每個(gè)消費(fèi)者至少消費(fèi)一次。
到此這篇關(guān)于RocketMQ的兩種消費(fèi)模式詳解的文章就介紹到這了,更多相關(guān)RocketMQ消費(fèi)模式內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
最新log4j2遠(yuǎn)程代碼執(zhí)行漏洞(附解決方法)
Apache?Log4j2?遠(yuǎn)程代碼執(zhí)行漏洞攻擊代碼,該漏洞利用無需特殊配置,經(jīng)多方驗(yàn)證,Apache?Struts2、Apache?Solr、Apache?Druid、Apache?Flink等均受影響,本文就介紹一下解決方法2021-12-12
Java?Post請求發(fā)送form-data表單參數(shù)詳細(xì)示例代碼
POST請求是一種常見的網(wǎng)絡(luò)通信操作,用于向服務(wù)器發(fā)送數(shù)據(jù),這種請求通常用于上傳文件或者提交包含大量數(shù)據(jù)的表單,這篇文章主要介紹了Java?Post請求發(fā)送form-data表單參數(shù)的相關(guān)資料,需要的朋友可以參考下2025-07-07
使用easyexcel導(dǎo)出的excel文件,使用poi讀取時(shí)異常處理方案
這篇文章主要介紹了使用easyexcel導(dǎo)出的excel文件,使用poi讀取時(shí)異常處理方案,具有很好的參考價(jià)值,希望對大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2023-12-12
Java使用EasyExcel動(dòng)態(tài)添加自增序號列
本文將介紹如何通過使用EasyExcel自定義攔截器實(shí)現(xiàn)在最終的Excel文件中新增一列自增的序號列,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2021-09-09
java設(shè)計(jì)模式學(xué)習(xí)之裝飾模式
這篇文章主要為大家詳細(xì)介紹了java設(shè)計(jì)模式學(xué)習(xí)之裝飾模式的相關(guān)資料,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2017-10-10

