最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

RocketMQ的兩種消費模式詳解

 更新時間:2023年10月11日 11:07:35   作者:碼奴生來只知道前進~  
這篇文章主要介紹了RocketMQ的兩種消費模式詳解,RocketMQ主要提供了兩種消費模式,集群消費以及廣播消費,我們只需要在定義消費者的時候通過setMessageModel(MessageModel.XXX),需要的朋友可以參考下

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)

// 設(shè)置消費模型,集群還是廣播,默認(rèn)為集群  CLUSTERING-集群,BROADCASTING-廣播
mqPushConsumer.setMessageModel(MessageModel.CLUSTERING);
mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);

方法就可以指定是集群還是廣播式消費,默認(rèn)是集群消費模式,即每個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類并設(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ā)送消息超時時間,默認(rèn)3000
        mqProducer.setSendMsgTimeout(3000);
        // 發(fā)送消息失敗重試次數(shù),默認(rèn)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(標(biāo)簽)和消息內(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ā)送結(jié)果
                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ā)送消息,關(guān)閉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類并設(shè)定消費者名稱
        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第一次啟動是從隊列頭部開始消費還是隊列尾部開始消費
        // 如果不是第一次啟動,那么按照上次消費的位置繼續(xù)消費
        mqPushConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
        // 設(shè)置消費模型,集群還是廣播,默認(rèn)為集群  CLUSTERING-集群   BROADCASTING-廣播
        mqPushConsumer.setMessageModel(MessageModel.CLUSTERING);
        // mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);
        // 消費者最小線程量
        mqPushConsumer.setConsumeThreadMin(5);
        // 消費者最大線程量
        mqPushConsumer.setConsumeThreadMax(10);
        // 設(shè)置一次消費消息的條數(shù),默認(rèn)是1
        mqPushConsumer.setConsumeMessageBatchMaxSize(1);
        // 訂閱一個或者多個Topic,以及Tag來過濾需要消費的消息,如果訂閱該主題下的所有tag,則使用*
        mqPushConsumer.subscribe("TOPIC_TEST", "*");
        // 注冊回調(diào)實現(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);*/
                // 標(biāo)記該消息已經(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類并設(shè)定消費者名稱
        DefaultMQPushConsumer mqPushConsumer = new DefaultMQPushConsumer(MqConstants.ConsumerGroup.CONSUMER_GROUP1);
        // 設(shè)置NameServer地址,如果是集群的話,使用分號;分隔開
        mqPushConsumer.setNamesrvAddr("127.0.0.1:9876");
        // 設(shè)置Consumer第一次啟動是從隊列頭部開始消費還是隊列尾部開始消費
        // 如果不是第一次啟動,那么按照上次消費的位置繼續(xù)消費
        mqPushConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
        // 設(shè)置消費模型,集群還是廣播,默認(rèn)為集群  CLUSTERING-集群   BROADCASTING-廣播
        // mqPushConsumer.setMessageModel(MessageModel.CLUSTERING);
        mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);
        // 消費者最小線程量
        mqPushConsumer.setConsumeThreadMin(5);
        // 消費者最大線程量
        mqPushConsumer.setConsumeThreadMax(10);
        // 設(shè)置一次消費消息的條數(shù),默認(rèn)是1
        mqPushConsumer.setConsumeMessageBatchMaxSize(1);
        // 訂閱一個或者多個Topic,以及Tag來過濾需要消費的消息,如果訂閱該主題下的所有tag,則使用*
        mqPushConsumer.subscribe("TOPIC_TEST", "*");
        // 注冊回調(diào)實現(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));
                }
                // 標(biāo)記該消息已經(jīng)被成功消費
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });
        // 啟動消費者實例
        mqPushConsumer.start();
        System.out.println("ConsumerB Started.");
    }
}

可以看到, 生產(chǎn)者發(fā)送了10條消息,ConsumerA與ConsumerB屬于同一個消費者組,集群消費模式下每個消費者攤分消費所有消息。注意,兩個消費者的ConsumerGroup組名需要一致,才算是同一個消費者組。

簡單總結(jié)一下:

1、在Rocket集群消費模式下,(訂閱)同一個主題(Topic)下的消息,對于不同的消費者組是一種“廣播形式”,即每個消費者組的都會消費消息。

2、在Rocket集群消費模式下,(訂閱)同一個主題(Topic)下的消息,對于相同的消費者組的消費者而言是一種集群模式,即同一個消費者組內(nèi)的所有消費者均分消息并消費。

 4、廣播消費

一條消息被多個 Consumer 消費,即使這些 Consumer 屬于同一個 Consumer Group,消息也會被 Consumer Group 中的每個 Consumer 都消費一次,廣播消費中的 Consumer Group 概念可以認(rèn)為在消息劃分方面無意 義。

使用方法:setMessageModel(MessageModel.BROADCASTING)

將前面的消費者A和消費者B的集群模式代碼設(shè)置為如下

mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);

重新啟動生成者和2個消費者

4.1 生產(chǎn)者數(shù)據(jù)

4.2 消費者A

4.3 消費者B

可以看到生產(chǎn)者發(fā)送了10條消息,ConsumerA與ConsumerB屬于同一個消費者組,廣播模式下每個消費者都會全量消費所有消息 。

  • 集群消費:任何一條消息只需要被消費者集群中任意一個消費者處理。
  • 廣播消費:每條消息被推送給消費者集群中的所有注冊消費者,保證消息被每個消費者至少消費一次。 

到此這篇關(guān)于RocketMQ的兩種消費模式詳解的文章就介紹到這了,更多相關(guān)RocketMQ消費模式內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 最新log4j2遠(yuǎn)程代碼執(zhí)行漏洞(附解決方法)

    最新log4j2遠(yuǎn)程代碼執(zhí)行漏洞(附解決方法)

    Apache?Log4j2?遠(yuǎn)程代碼執(zhí)行漏洞攻擊代碼,該漏洞利用無需特殊配置,經(jīng)多方驗證,Apache?Struts2、Apache?Solr、Apache?Druid、Apache?Flink等均受影響,本文就介紹一下解決方法
    2021-12-12
  • Java?Post請求發(fā)送form-data表單參數(shù)詳細(xì)示例代碼

    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
  • JFinal極速開發(fā)框架使用筆記分享

    JFinal極速開發(fā)框架使用筆記分享

    下面小編就為大家分享一篇JFinal極速開發(fā)框架使用筆記,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2018-01-01
  • 使用easyexcel導(dǎo)出的excel文件,使用poi讀取時異常處理方案

    使用easyexcel導(dǎo)出的excel文件,使用poi讀取時異常處理方案

    這篇文章主要介紹了使用easyexcel導(dǎo)出的excel文件,使用poi讀取時異常處理方案,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2023-12-12
  • SpringBoot整合Lombok及常見問題解決

    SpringBoot整合Lombok及常見問題解決

    本文主要介紹了SpringBoot整合Lombok及常見問題解決,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2022-04-04
  • javap命令的使用技巧

    javap命令的使用技巧

    本篇文章給大家分享了關(guān)于JAVA中關(guān)于javap命令的使用技巧以及相關(guān)代碼分享,有需要的朋友參考學(xué)習(xí)下。
    2018-05-05
  • Java使用EasyExcel動態(tài)添加自增序號列

    Java使用EasyExcel動態(tài)添加自增序號列

    本文將介紹如何通過使用EasyExcel自定義攔截器實現(xiàn)在最終的Excel文件中新增一列自增的序號列,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2021-09-09
  • java設(shè)計模式學(xué)習(xí)之裝飾模式

    java設(shè)計模式學(xué)習(xí)之裝飾模式

    這篇文章主要為大家詳細(xì)介紹了java設(shè)計模式學(xué)習(xí)之裝飾模式的相關(guān)資料,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2017-10-10
  • Spring Boot常見外部配置文件方式詳析

    Spring Boot常見外部配置文件方式詳析

    這篇文章主要給大家介紹了關(guān)于Spring Boot常見外部配置文件方式的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者使用Spring Boot具有一定的參考學(xué)習(xí)價值,需要的朋友們下面來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-07-07
  • idea編寫java程序詳細(xì)圖文步驟

    idea編寫java程序詳細(xì)圖文步驟

    這篇文章主要給大家介紹了關(guān)于idea編寫java程序的詳細(xì)圖文步驟,IDEA是用于Java語言開發(fā)的集成環(huán)境,它是業(yè)界公認(rèn)的目前用于Java程序開發(fā)最好的工具,文中通過圖文介紹的非常詳細(xì),需要的朋友可以參考下
    2023-09-09

最新評論

萨迦县| 西城区| 崇礼县| 寿宁县| 兰西县| 新竹县| 宁夏| 石狮市| 万全县| 拉孜县| 大洼县| 新兴县| 龙口市| 祁东县| 滦平县| 迁安市| 湖南省| 黄大仙区| 化州市| 塔河县| 永新县| 乐业县| 长白| 秭归县| 文化| 忻州市| 科技| 东兴市| 澄城县| 赤城县| 南充市| 潞西市| 蓝山县| 澄迈县| 阿尔山市| 商南县| 罗城| 南阳市| 漳平市| 会泽县| 宜兰县|