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

SpringBoot整合RocketMQ批量發(fā)送消息的實現(xiàn)代碼

 更新時間:2024年04月01日 10:41:26   作者:嘉禾嘉寧papa  
這篇文章主要介紹了SpringBoot整合RocketMQ批量發(fā)送消息的實現(xiàn),文中通過代碼示例講解的非常詳細(xì),對大家的學(xué)習(xí)或工作有一定的幫助,需要的朋友可以參考下

一、簡介

今天我們講講如何批量發(fā)送消息,主要還是使用方法RocketMQTemplatesyncSend方法。

1.1、特點

批量發(fā)送和單條發(fā)送消息的主要區(qū)別有以下幾點:

  • 網(wǎng)絡(luò)開銷 發(fā)送單條消息時,每個消息都需要單獨建立網(wǎng)絡(luò)連接、發(fā)送數(shù)據(jù)包、等待響應(yīng)等,網(wǎng)絡(luò)開銷較大。批量發(fā)送可以將多條消息打包在一起發(fā)送,減少網(wǎng)絡(luò)連接建立的次數(shù),降低網(wǎng)絡(luò)開銷
  • 吞吐量 由于批量發(fā)送減少了網(wǎng)絡(luò)開銷,所以可以在單位時間內(nèi)發(fā)送更多的消息,提高了吞吐量。在高并發(fā)高流量場景下,批量發(fā)送能夠發(fā)揮更好的性能
  • 消息順序 單條發(fā)送消息的順序是有序的,后發(fā)送的在隊列中排在前發(fā)送的后面。而對于批量發(fā)送,一個批次內(nèi)的消息順序是固定的,但不同批次之間的消息順序是無序的,會按照到達(dá)順序存儲在隊列中。如果需要嚴(yán)格消息順序,單條發(fā)送更合適
  • 消息重試 如果批量發(fā)送的一個批次中有部分消息發(fā)送失敗,需重發(fā)整個批次,沒有選擇重發(fā)其中部分消息的功能(涉及冪等性問題)。單條發(fā)送失敗時只需重發(fā)該單條消息
  • 編程復(fù)雜度 批量發(fā)送需要構(gòu)造MessageBatch或Message列表對象,編程略微復(fù)雜些。單條發(fā)送只需構(gòu)造單個Message對象

二、Maven依賴

pom.xml

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <parent>
        <artifactId>rocketmq</artifactId>
        <groupId>com.alian</groupId>
        <version>1.0.0-SNAPSHOT</version>
    </parent>
    <modelVersion>4.0.0</modelVersion>

    <artifactId>06-send-batched-message</artifactId>

    <properties>
        <maven.compiler.source>8</maven.compiler.source>
        <maven.compiler.target>8</maven.compiler.target>
    </properties>

    <dependencies>
        <dependency>
            <groupId>com.alian</groupId>
            <artifactId>common-rocketmq-dto</artifactId>
            <version>1.0.0-SNAPSHOT</version>
        </dependency>
    </dependencies>

</project>

父工程已經(jīng)在我上一篇文章里,通用公共包也在我上一篇文章里有說明,包括消費者。具體參考:SpringBoot整合RocketMQ實現(xiàn)發(fā)送同步消息_java_腳本之家 (jb51.net)

三、application配置

application.properties

server.port=8005

# rocketmq地址
rocketmq.name-server=192.168.0.234:9876
# 默認(rèn)的生產(chǎn)者組
rocketmq.producer.group=batched_group
# 發(fā)送同步消息超時時間
rocketmq.producer.send-message-timeout=3000
# 用于設(shè)置在消息發(fā)送失敗后,生產(chǎn)者是否嘗試切換到下一個服務(wù)器。設(shè)置為 true 表示啟用,在發(fā)送失敗時嘗試切換到下一個服務(wù)器
rocketmq.producer.retry-next-server=true
# 用于指定消息發(fā)送失敗時的重試次數(shù)
rocketmq.producer.retry-times-when-send-failed=3
# 設(shè)置消息壓縮的閾值,為0表示禁用消息體的壓縮
rocketmq.producer.compress-message-body-threshold=0

四、批量發(fā)送

在 RocketMQ 中,RocketMQTemplatesyncSend方法,它允許你批量發(fā)送同步消息,主要參數(shù):

  • topic:(普通消息都發(fā)送到topic=string_message_topic
  • Collection<T>:消息集合

測試類都引入依賴

	@Autowired
    private RocketMQTemplate rocketMQTemplate;

4.1、同步消息

    @Test
    public void syncSendBatchStringMessagesWithBuilder() {
        String topic = "string_message_topic";
        String message = "超級喜歡Golang語言:";
        List<Message<String>> messageList = new ArrayList<>();
        for (int i = 0; i < 10; i++) {
            Message<String> rocketMessage = MessageBuilder.withPayload(message + i)
                    // 設(shè)置消息類型
                    .setHeader(MessageHeaders.CONTENT_TYPE, "text/plain")
                    .build();
            // 加入到列表
            messageList.add(rocketMessage);
        }
        // 使用syncSend發(fā)送批量消息
        SendResult sendResult = rocketMQTemplate.syncSend(topic, messageList);
        log.info("同步批量發(fā)送普通消息結(jié)果:{}",sendResult);
    }

運行結(jié)果:

[_GROUP_STRING_1] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 同步批量發(fā)送普通消息:0
[_GROUP_STRING_2] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 同步批量發(fā)送普通消息:1
[_GROUP_STRING_4] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 同步批量發(fā)送普通消息:3
[_GROUP_STRING_5] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 同步批量發(fā)送普通消息:4
[_GROUP_STRING_3] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 同步批量發(fā)送普通消息:2
[_GROUP_STRING_6] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 同步批量發(fā)送普通消息:5
[_GROUP_STRING_7] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 同步批量發(fā)送普通消息:6
[_GROUP_STRING_8] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 同步批量發(fā)送普通消息:7
[_GROUP_STRING_9] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 同步批量發(fā)送普通消息:8
[GROUP_STRING_10] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 同步批量發(fā)送普通消息:9

4.2、異步消息

    @Test
    public void asyncSendBatchStringMessageWithBuilder() {
        String topic = "string_message_topic";
        String message = "Alian超級喜歡Golang語言:";
        List<Message<String>> messageList = new ArrayList<>();
        for (int i = 0; i < 10; i++) {
            Message<String> rocketMessage = MessageBuilder.withPayload(message + i)
                    // 設(shè)置消息類型
                    .setHeader(MessageHeaders.CONTENT_TYPE, "text/plain")
                    .build();
            // 加入到列表
            messageList.add(rocketMessage);
        }
        // 使用asyncSend發(fā)送批量消息
        rocketMQTemplate.asyncSend(topic, messageList, new SendCallback() {
            @Override
            public void onSuccess(SendResult sendResult) {
                // 異步發(fā)送成功的回調(diào)邏輯
                log.info("異步批量發(fā)送普通消息成功: " + sendResult);
            }

            @Override
            public void onException(Throwable e) {
                // 異步發(fā)送失敗的回調(diào)邏輯
                log.info("異步批量發(fā)送普通消息失敗: " + e.getMessage());
            }
        });
    }

運行結(jié)果:

[_GROUP_STRING_1] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 異步批量發(fā)送普通消息:0
[_GROUP_STRING_8] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 異步批量發(fā)送普通消息:1
[_GROUP_STRING_3] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 異步批量發(fā)送普通消息:7
[_GROUP_STRING_6] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 異步批量發(fā)送普通消息:4
[_GROUP_STRING_9] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 異步批量發(fā)送普通消息:2
[_GROUP_STRING_5] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 異步批量發(fā)送普通消息:6
[GROUP_STRING_10] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 異步批量發(fā)送普通消息:3
[_GROUP_STRING_2] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 異步批量發(fā)送普通消息:8
[_GROUP_STRING_4] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 異步批量發(fā)送普通消息:9
[_GROUP_STRING_7] c.a.concurrent.StringMessageConsumer     : 字符串消費者接收到的消息: 異步批量發(fā)送普通消息:5

4.3、順序消息

在 RocketMQ 中,RocketMQTemplatesyncSendOrderly方法,它允許你批量發(fā)送同步消息,主要參數(shù):

  • topic:(和之前有區(qū)別,普通消息都發(fā)送到topic=ordered_string_message_topic
  • Collection<T>:消息集合
  • hashKey:通過hashKey發(fā)送到同一個隊列
    @Test
    public void syncSendBatchOrderlyStringMessagesWithBuilder() {
        String topic = "ordered_string_message_topic";
        String message = "同步批量發(fā)送順序消息,超級喜歡Go語言:";
        List<Message<String>> messageList = new ArrayList<>();
        for (int i = 0; i < 10; i++) {
            Message<String> rocketMessage = MessageBuilder.withPayload(message + i)
                    // 設(shè)置消息類型
                    .setHeader(MessageHeaders.CONTENT_TYPE, "text/plain")
                    .build();
            // 加入到列表
            messageList.add(rocketMessage);
        }
        // 使用syncSendOrderly發(fā)送批量順序消息,消費者線程設(shè)置為1
        SendResult sendResult = rocketMQTemplate.syncSendOrderly(topic, messageList, "alian_sync_ordered");
        log.info("批量發(fā)送順序消息發(fā)送結(jié)果:{}",sendResult);
    }

運行結(jié)果:

[GROUP_STRING_10] com.alian.ordered.StringMessageConsumer  : 字符串消費者接收到的消息: 同步批量發(fā)送順序消息,超級喜歡Go語言:0
[GROUP_STRING_10] com.alian.ordered.StringMessageConsumer  : 字符串消費者接收到的消息: 同步批量發(fā)送順序消息,超級喜歡Go語言:1
[GROUP_STRING_10] com.alian.ordered.StringMessageConsumer  : 字符串消費者接收到的消息: 同步批量發(fā)送順序消息,超級喜歡Go語言:2
[GROUP_STRING_10] com.alian.ordered.StringMessageConsumer  : 字符串消費者接收到的消息: 同步批量發(fā)送順序消息,超級喜歡Go語言:3
[GROUP_STRING_10] com.alian.ordered.StringMessageConsumer  : 字符串消費者接收到的消息: 同步批量發(fā)送順序消息,超級喜歡Go語言:4
[GROUP_STRING_10] com.alian.ordered.StringMessageConsumer  : 字符串消費者接收到的消息: 同步批量發(fā)送順序消息,超級喜歡Go語言:5
[GROUP_STRING_10] com.alian.ordered.StringMessageConsumer  : 字符串消費者接收到的消息: 同步批量發(fā)送順序消息,超級喜歡Go語言:6
[GROUP_STRING_10] com.alian.ordered.StringMessageConsumer  : 字符串消費者接收到的消息: 同步批量發(fā)送順序消息,超級喜歡Go語言:7
[GROUP_STRING_10] com.alian.ordered.StringMessageConsumer  : 字符串消費者接收到的消息: 同步批量發(fā)送順序消息,超級喜歡Go語言:8
[GROUP_STRING_10] com.alian.ordered.StringMessageConsumer  : 字符串消費者接收到的消息: 同步批量發(fā)送順序消息,超級喜歡Go語言:9

所以我之前說批量發(fā)送消息的topic不一樣,因為

@Slf4j
@Service
@RocketMQMessageListener(topic = "ordered_string_message_topic", 
						consumerGroup = "ORDERED_GROUP_STRING", 
						consumeMode = ConsumeMode.ORDERLY)
public class StringMessageConsumer implements RocketMQListener<String> {

    @Override
    public void onMessage(String message) {
        log.info("字符串消費者接收到的消息: {}", message);
        // 處理消息的業(yè)務(wù)邏輯
    }
}

順序消息要順序消費,也就是每次是一個線程去消費,相當(dāng)于單線程,也就有序了。關(guān)鍵就是配置了:consumeMode = ConsumeMode.ORDERLY

當(dāng)然,我們也可以把消費者線程數(shù)設(shè)置為 consumeThreadNumber = 1,也就是單線程消費了,從而確保了消息的順序消費(指單實例):

@RocketMQMessageListener(topic = "ordered_string_message_topic", 
						consumerGroup = "CONCURRENT_GROUP_STRING",
						consumeThreadNumber = 1)
public class StringMessageConsumer implements RocketMQListener<String> {

    @Override
    public void onMessage(String message) {
        log.info("字符串消費者接收到的消息: {}", message);
        // 處理消息的業(yè)務(wù)邏輯
    }
}

4.4、關(guān)于異步批量發(fā)送

有可能你會寫下面的異步批量發(fā)送順序消息

	@Test
    public void asyncSendBatchOrderlyStringMessageWithBuilder2() {
        String topic = "ordered_string_message_topic";
        String message = "Alian超級喜歡Golang語言:";
        List<String> messageList = new ArrayList<>();
        for (int i = 0; i < 10; i++) {
            // 加入到列表
            messageList.add(message + i);
        }
        // 使用 asyncSendOrderly 發(fā)送批量消息
        rocketMQTemplate.asyncSendOrderly(topic, messageList, "alian_async_ordered", new SendCallback() {
            @Override
            public void onSuccess(SendResult sendResult) {
                // 異步發(fā)送成功的回調(diào)邏輯
                log.info("異步消息發(fā)送字符串消息成功: " + sendResult);
            }

            @Override
            public void onException(Throwable e) {
                // 異步發(fā)送失敗的回調(diào)邏輯
                log.info("異步消息發(fā)送字符串消息失敗: " + e.getMessage());
            }
        });
    }

其實這個是不對的,最終的結(jié)果是一個把你這里的messageList,當(dāng)做了一個消息列表接收了,如下結(jié)果:

[GROUP_STRING_18] com.alian.ordered.StringMessageConsumer  : 字符串消費者接收到的消息: ["Alian超級喜歡Golang語言:0","Alian超級喜歡Golang語言:1","Alian超級喜歡Golang語言:2","Alian超級喜歡Golang語言:3","Alian超級喜歡Golang語言:4","Alian超級喜歡Golang語言:5","Alian超級喜歡Golang語言:6","Alian超級喜歡Golang語言:7","Alian超級喜歡Golang語言:8","Alian超級喜歡Golang語言:9"]

RocketMQ對于單條消息和批量消息在隊列中是如何被處理的?

  • 對于單條發(fā)送的消息,RocketMQ會按照隊列中的順序,將每條消息分發(fā)給一個消費者線程。因此,即使有多個消費者線程,由于每條消息都被單獨處理,消費的順序仍然會與發(fā)送的順序一致。
  • 對于批量發(fā)送的消息,情況就有所不同。批量消息是作為一個整體發(fā)送的,因此在隊列中,它們被視為一個單獨的實體。當(dāng)RocketMQ從隊列中取出批量消息時,它會將整個批量消息作為一個整體分發(fā)給一個消費者線程。如果有多個消費者線程,由于操作系統(tǒng)的線程調(diào)度策略,處理批量消息的線程可能會在處理消息的過程中被調(diào)度出去,從而允許其他線程處理后面的消息。這樣就可能導(dǎo)致消費的順序與發(fā)送的順序不一致。

4.5、結(jié)論

為此我測試了多次,得到結(jié)論:

  • 單條發(fā)送消息到同一個隊列,使用多個消費線程消費該隊列,由于消息本身是有序的,所以消費順序也是有序的
  • 單批次批量發(fā)送消息到同一個隊列,使用單個消費線程消費該隊列,由于消費線程是單一的,所以消費順序也是有序的
  • 單批次批量發(fā)送消息到同一個隊列,使用多個消費線程消費時,消費順序就不是有序的了

五、其他

既然知道批量消息是作為一個整體的,那么肯定就會對消息大小有限制,在 Apache RocketMQ 中,批量消息的大小默認(rèn)限制是4MB。這意味著,你不能發(fā)送總大小超過4MB的批量消息。

如果你想修改這個限制,你需要修改RocketMQ的配置。具體的修改方法如下:

  • 找到RocketMQ的配置文件broker.conf,這個文件通常位于RocketMQ安裝目錄的conf目錄下。
  • 在broker.conf文件中,找到maxMessageSize這個配置項。這個配置項決定了批量消息的最大大小。
  • 修改maxMessageSize的值為你想要的大小。注意,這個值是以字節(jié)為單位的,所以如果你想設(shè)置批量消息的最大大小為8MB,你應(yīng)該設(shè)置maxMessageSize=8388608。
  • 保存并關(guān)閉broker.conf文件。
  • 重啟RocketMQ的Broker服務(wù),以使新的配置生效。

雖然你可以通過修改配置來增加批量消息的最大大小,但是你應(yīng)該謹(jǐn)慎地考慮這個決定。增加批量消息的最大大小可能會增加Broker的內(nèi)存使用量,并可能影響到消息的發(fā)送和接收性能。因此,在修改這個配置之前,你應(yīng)該先考慮你的應(yīng)用的需求和Broker的性能。

因為優(yōu)先的是@RocketMQMessageListener 注解中設(shè)置 consumerGroup 和messageModel 參數(shù)。

以上就是SpringBoot整合RocketMQ批量發(fā)送消息的實現(xiàn)的詳細(xì)內(nèi)容,更多關(guān)于SpringBoot RocketMQ批量發(fā)送消息的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • JDK1.8中ConcurrentHashMap中computeIfAbsent死循環(huán)bug問題

    JDK1.8中ConcurrentHashMap中computeIfAbsent死循環(huán)bug問題

    這篇文章主要介紹了JDK1.8中ConcurrentHashMap中computeIfAbsent死循環(huán)bug,本文給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2020-08-08
  • SpringBoot中5種動態(tài)代理的實現(xiàn)方案

    SpringBoot中5種動態(tài)代理的實現(xiàn)方案

    在SpringBoot應(yīng)用中,動態(tài)代理被廣泛用于實現(xiàn)事務(wù)管理、緩存、安全控制、日志記錄等橫切關(guān)注點,下面小編為大家介紹一下SpringBoot中5種動態(tài)代理的實現(xiàn)方案吧
    2025-04-04
  • SpringBoot啟動yaml報錯的解決

    SpringBoot啟動yaml報錯的解決

    這篇文章主要介紹了SpringBoot啟動yaml報錯的解決方案,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-08-08
  • SpringBoot淺析安全管理之基于數(shù)據(jù)庫認(rèn)證

    SpringBoot淺析安全管理之基于數(shù)據(jù)庫認(rèn)證

    在真實的項目中,用戶的基本信息以及角色等都存儲在數(shù)據(jù)庫中,因此需要從數(shù)據(jù)庫中獲取數(shù)據(jù)進(jìn)行認(rèn)證和授權(quán)
    2022-08-08
  • 淺析Java如何防范SQL注入攻擊

    淺析Java如何防范SQL注入攻擊

    這篇文章主要為大家詳細(xì)介紹了SQL注入漏洞的原理,表現(xiàn)形式以及如何防范這種攻擊,文中的示例代碼講解詳細(xì),希望可以幫助開發(fā)者構(gòu)建更安全的Java應(yīng)用
    2025-05-05
  • 本地JDK多版本快速切換方案

    本地JDK多版本快速切換方案

    本文將詳細(xì)介紹如何在同一臺機器上安裝和配置多個版本的 JDK(JDK 8、JDK 17 和 JDK 21),并且使用綠色版(即無需安裝程序,直接解壓即可使用),通過這種方式,您可以在不同的項目中靈活選擇所需的 JDK 版本,需要的朋友可以參考下
    2025-03-03
  • RabbitMQ?Stream插件使用案例代碼

    RabbitMQ?Stream插件使用案例代碼

    這篇文章主要介紹了RabbitMQ?Stream插件使用案例代碼,2.4版為RabbitMQ流插件引入了對RabbitMQStream插件Java客戶端的初始支持,需要的朋友可以參考下
    2024-04-04
  • Java編程之如何通過JSP實現(xiàn)頭像自定義上傳

    Java編程之如何通過JSP實現(xiàn)頭像自定義上傳

    之前做這個頭像上傳功能還是花了好多時間的,今天我將我的代碼分享給大家,下面這篇文章主要給大家介紹了關(guān)于Java編程之如何通過JSP實現(xiàn)頭像自定義上傳的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2022-12-12
  • java中實現(xiàn)控制臺打印sql語句方式

    java中實現(xiàn)控制臺打印sql語句方式

    這篇文章主要介紹了java中實現(xiàn)控制臺打印sql語句方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2023-06-06
  • 使用Java和Redis實現(xiàn)高效的短信防轟炸方案

    使用Java和Redis實現(xiàn)高效的短信防轟炸方案

    在當(dāng)今互聯(lián)網(wǎng)應(yīng)用中,短信驗證碼已成為身份驗證的重要手段,然而,這也帶來了"短信轟炸"的安全風(fēng)險?-?惡意用戶利用程序自動化發(fā)送大量短信請求,導(dǎo)致用戶被騷擾和企業(yè)短信成本激增,本文將詳細(xì)介紹如何使用Java和Redis實現(xiàn)高效的短信防轟炸解決方案,需要的朋友可以參考下
    2025-04-04

最新評論

察哈| 连云港市| 墨竹工卡县| 荣成市| 湖南省| 龙泉市| 呼和浩特市| 东山县| 新昌县| 手游| 永靖县| 邮箱| 诸暨市| 阿合奇县| 林芝县| 卓资县| 库伦旗| 泸溪县| 赤城县| 黔东| 云阳县| 吉木萨尔县| 高陵县| 宕昌县| 柞水县| 巴南区| 邯郸县| 郑州市| 成安县| 凤凰县| 利辛县| 湘西| 锦州市| 潜江市| 九江县| 涿州市| 垫江县| 灵寿县| 桃园市| 枣庄市| 平果县|