RocketMQ的兩種消費模式詳解
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、消費模式
RocketMQ主要提供了兩種消費模式:集群消費以及廣播消費。我們只需要在定義消費者的時候通過setMessageModel(MessageModel.XXX)
// 設置消費模型,集群還是廣播,默認為集群 CLUSTERING-集群,BROADCASTING-廣播 mqPushConsumer.setMessageModel(MessageModel.CLUSTERING); mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);
方法就可以指定是集群還是廣播式消費,默認是集群消費模式,即每個Consumer Group中的Consumer均攤所有的消息。
3、集群消費
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類并設定生產(chǎn)者名稱 DefaultMQProducer mqProducer = new DefaultMQProducer("producer-group-test"); // 設置NameServer地址,如果是集群的話,使用分號;分隔開 mqProducer.setNamesrvAddr("127.0.0.1:9876"); // 消息最大長度 默認4M mqProducer.setMaxMessageSize(4096); // 發(fā)送消息超時時間,默認3000 mqProducer.setSendMsgTimeout(3000); // 發(fā)送消息失敗重試次數(shù),默認2 mqProducer.setRetryTimesWhenSendAsyncFailed(2); // 啟動消息生產(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(標簽)和消息內(nèi)容 Message message = new Message("TOPIC_TEST", "", msg.getBytes(RemotingHelper.DEFAULT_CHARSET)); // 發(fā)送同步消息到一個Broker,可以通過sendResult返回消息是否成功送達 SendResult sendResult = mqProducer.send(message); // mqProducer.sendOneway(message); // 消息id /*System.out.println(sendResult.getMsgId()); // 隊列信息 System.out.println(sendResult.getMessageQueue()); // 發(fā)送結果 System.out.println(sendResult.getSendStatus()); // 下一個要消費的消息的偏移量 System.out.println(sendResult.getOffsetMsgId()); // 隊列消息偏移量 System.out.println(sendResult.getQueueOffset());*/ System.out.println(sendResult); } } catch (Exception e) { e.printStackTrace(); System.out.println("生產(chǎn)消息異常!"); } // 如果不再發(fā)送消息,關閉Producer實例 mqProducer.shutdown(); } }
User用戶測試實體類
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)建時間 */ @JsonFormat(pattern = DateUtil.DATETIME_FORMAT) private Date createTime; }
3.2 消費者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 消費者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類并設定消費者名稱 DefaultMQPushConsumer mqPushConsumer = new DefaultMQPushConsumer(MqConstants.ConsumerGroup.CONSUMER_GROUP1); // DefaultMQPullConsumer pullConsumer = new DefaultMQPullConsumer(MqConstants.ConsumerGroup.CONSUMER_GROUP1); // 設置NameServer地址,如果是集群的話,使用分號;分隔開 mqPushConsumer.setNamesrvAddr("127.0.0.1:9876"); // pullConsumer.setNamesrvAddr("127.0.0.1:9876"); // 設置Consumer第一次啟動是從隊列頭部開始消費還是隊列尾部開始消費 // 如果不是第一次啟動,那么按照上次消費的位置繼續(xù)消費 mqPushConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); // 設置消費模型,集群還是廣播,默認為集群 CLUSTERING-集群 BROADCASTING-廣播 mqPushConsumer.setMessageModel(MessageModel.CLUSTERING); // mqPushConsumer.setMessageModel(MessageModel.BROADCASTING); // 消費者最小線程量 mqPushConsumer.setConsumeThreadMin(5); // 消費者最大線程量 mqPushConsumer.setConsumeThreadMax(10); // 設置一次消費消息的條數(shù),默認是1 mqPushConsumer.setConsumeMessageBatchMaxSize(1); // 訂閱一個或者多個Topic,以及Tag來過濾需要消費的消息,如果訂閱該主題下的所有tag,則使用* mqPushConsumer.subscribe("TOPIC_TEST", "*"); // 注冊回調實現(xiàn)類來處理從broker拉取回來的消息 mqPushConsumer.registerMessageListener(new MessageListenerConcurrently() { // 監(jiān)聽類實現(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("消費者A接收到消息:" + "===== " + msgBody); /*User user = JSON.parseObject(msgBody, User.class); System.out.println("消費者A接收到消息:" + "===== " + user.getId());*/ } /*MessageExt messageExt = msgList.get(0); String body = new String(messageExt.getBody(), StandardCharsets.UTF_8); System.out.println("消費者A接收到消息: " + messageExt.toString() + "---消息內(nèi)容為:" + body);*/ // 標記該消息已經(jīng)被成功消費 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); // 啟動消費者實例 mqPushConsumer.start(); System.out.println("ConsumerA Started."); } }
3.3 消費者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 消費者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類并設定消費者名稱 DefaultMQPushConsumer mqPushConsumer = new DefaultMQPushConsumer(MqConstants.ConsumerGroup.CONSUMER_GROUP1); // 設置NameServer地址,如果是集群的話,使用分號;分隔開 mqPushConsumer.setNamesrvAddr("127.0.0.1:9876"); // 設置Consumer第一次啟動是從隊列頭部開始消費還是隊列尾部開始消費 // 如果不是第一次啟動,那么按照上次消費的位置繼續(xù)消費 mqPushConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); // 設置消費模型,集群還是廣播,默認為集群 CLUSTERING-集群 BROADCASTING-廣播 // mqPushConsumer.setMessageModel(MessageModel.CLUSTERING); mqPushConsumer.setMessageModel(MessageModel.BROADCASTING); // 消費者最小線程量 mqPushConsumer.setConsumeThreadMin(5); // 消費者最大線程量 mqPushConsumer.setConsumeThreadMax(10); // 設置一次消費消息的條數(shù),默認是1 mqPushConsumer.setConsumeMessageBatchMaxSize(1); // 訂閱一個或者多個Topic,以及Tag來過濾需要消費的消息,如果訂閱該主題下的所有tag,則使用* mqPushConsumer.subscribe("TOPIC_TEST", "*"); // 注冊回調實現(xiàn)類來處理從broker拉取回來的消息 mqPushConsumer.registerMessageListener(new MessageListenerConcurrently() { // 監(jiān)聽類實現(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("消費者B接收到消息: " + messageExt.toString() + "---消息內(nèi)容為:" + body);*/ for (MessageExt msg : msgList) { System.out.println("消費者B接收到消息:" + "===== " + new String(msg.getBody())); // String msgBody = new String(msg.getBody(), Charset.forName(RemotingHelper.DEFAULT_CHARSET)); } // 標記該消息已經(jīng)被成功消費 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); // 啟動消費者實例 mqPushConsumer.start(); System.out.println("ConsumerB Started."); } }
可以看到, 生產(chǎn)者發(fā)送了10條消息,ConsumerA與ConsumerB屬于同一個消費者組,集群消費模式下每個消費者攤分消費所有消息。注意,兩個消費者的ConsumerGroup組名需要一致,才算是同一個消費者組。
簡單總結一下:
1、在Rocket集群消費模式下,(訂閱)同一個主題(Topic)下的消息,對于不同的消費者組是一種“廣播形式”,即每個消費者組的都會消費消息。
2、在Rocket集群消費模式下,(訂閱)同一個主題(Topic)下的消息,對于相同的消費者組的消費者而言是一種集群模式,即同一個消費者組內(nèi)的所有消費者均分消息并消費。
4、廣播消費
一條消息被多個 Consumer 消費,即使這些 Consumer 屬于同一個 Consumer Group,消息也會被 Consumer Group 中的每個 Consumer 都消費一次,廣播消費中的 Consumer Group 概念可以認為在消息劃分方面無意 義。
使用方法:setMessageModel(MessageModel.BROADCASTING)
將前面的消費者A和消費者B的集群模式代碼設置為如下
mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);
重新啟動生成者和2個消費者
4.1 生產(chǎn)者數(shù)據(jù)
4.2 消費者A
4.3 消費者B
可以看到生產(chǎn)者發(fā)送了10條消息,ConsumerA與ConsumerB屬于同一個消費者組,廣播模式下每個消費者都會全量消費所有消息 。
- 集群消費:任何一條消息只需要被消費者集群中任意一個消費者處理。
- 廣播消費:每條消息被推送給消費者集群中的所有注冊消費者,保證消息被每個消費者至少消費一次。
到此這篇關于RocketMQ的兩種消費模式詳解的文章就介紹到這了,更多相關RocketMQ消費模式內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!
相關文章
Spring Boot實現(xiàn)跨域訪問實現(xiàn)代碼
本文通過實例代碼給大家介紹了Spring Boot實現(xiàn)跨域訪問的知識,然后在文中給大家介紹了spring boot 服務器端設置允許跨域訪問 的方法,感興趣的朋友一起看看吧2017-07-07Java練手小項目實現(xiàn)一個項目管理系統(tǒng)
讀萬卷書不如行萬里路,只學書上的理論是遠遠不夠的,只有在實戰(zhàn)中才能獲得能力的提升,本篇文章手把手帶你用Java實現(xiàn)一個項目管理系統(tǒng),大家可以在過程中查缺補漏,提升水平2021-10-10MyBatis-Plus 分頁查詢以及自定義sql分頁的實現(xiàn)
這篇文章主要介紹了MyBatis-Plus 分頁查詢以及自定義sql分頁的實現(xiàn),文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧2020-08-08如何使用Spring AOP的通知類型及創(chuàng)建通知
這篇文章主要給大家介紹了關于如何使用Spring AOP的通知類型及創(chuàng)建通知的相關資料,文中通過示例代碼介紹的非常詳細,對大家學習或者使用Spring AOP具有一定的參考學習價值,需要的朋友們下面來一起學習學習吧2019-12-12Java中統(tǒng)計字符個數(shù)以及反序非相同字符的方法詳解
本篇文章是對Java中統(tǒng)計字符個數(shù)以及反序非相同字符的方法進行了詳細的分析介紹,需要的朋友參考下2013-05-05