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

Rocketmq消息批量發(fā)送&消息批量消費方式

 更新時間:2026年05月16日 09:27:25   作者:拽著尾巴的魚兒  
文章介紹了RocketMQ中消息的批量發(fā)送和消費機制,包括批量發(fā)送的優(yōu)點、限制及注意事項,以及消息消費模式的推、拉兩種方式,重點闡述了DefaultLitePullConsumer和DefaultMQPushConsumer的使用方法和注意事項,并提供了代碼示例

前言:批量發(fā)送和消費消息在一定程度上可以提高吞吐量,減少帶寬,那么Rocketmq 中的消息怎么進(jìn)行批量的發(fā)送和批量的消費呢;

1 消息的批量發(fā)送

1.1 批量發(fā)送的優(yōu)點以及實現(xiàn)

批量發(fā)送消息可以提高 RocketMQ 的生產(chǎn)者性能和吞吐量。由于批量發(fā)送消息可以減少網(wǎng)絡(luò) I/O 操作和降低消息發(fā)送延遲,因此它在以下情況下特別有用:

  • 發(fā)送大量小型消息時
  • 需要降低消息發(fā)送延遲時
  • 需要提高生產(chǎn)者性能時

但是,批量發(fā)送消息也存在一些注意事項,需要注意以下幾點:

  • 消息列表的大小不能超過 broker 設(shè)置的最大消息大小。
  • 消息列表的大小不能超過生產(chǎn)者設(shè)置的 maxMessageSize 參數(shù),此參數(shù)默認(rèn)為 4MB。
  • 批量發(fā)送消息不支持消息事務(wù)。
  • 如果的代碼在發(fā)送消息列表時發(fā)生異常,則可能會發(fā)生部分消息發(fā)送成功,部分消息發(fā)送失敗的情況。如果要確保所有消息都已成功發(fā)送,則需要增加錯誤處理邏輯和消息重試機制。

批量發(fā)送消息是一種提高 RocketMQ 生產(chǎn)者性能和吞吐量的好方法,但需要注意消息列表大小和錯誤處理機制,以確保生產(chǎn)者的可靠性和穩(wěn)定性。

public class SimpleBatchProducer {

    public static void main(String[] args) throws Exception {
        DefaultMQProducer producer = new DefaultMQProducer("BatchProducerGroupName");
        producer.start();

        //If you just send messages of no more than 1MiB at a time, it is easy to use batch
        //Messages of the same batch should have: same topic, same waitStoreMsgOK and no schedule support
        String topic = "BatchTest";
        List<Message> messages = new ArrayList<>();
        messages.add(new Message(topic, "Tag", "OrderID001", "Hello world 0".getBytes()));
        messages.add(new Message(topic, "Tag", "OrderID002", "Hello world 1".getBytes()));
        messages.add(new Message(topic, "Tag", "OrderID003", "Hello world 2".getBytes()));

        producer.send(messages);
    }
}

1.2 批量發(fā)送消息為什么要限制maxMessageSize 

消息列表的大小不能超過生產(chǎn)者設(shè)置的 maxMessageSize 參數(shù),主要是為了避免消息發(fā)送延遲和消息過大導(dǎo)致 broker 出現(xiàn)性能問題。如果嘗試發(fā)送大于 maxMessageSize 的消息,RocketMQ 會拋出 MessageTooLargeException 異常,并且消息不會被發(fā)送到 broker。

如果開發(fā)者在開發(fā)時遇到了消息列表大小超過 maxMessageSize 的情況,可以考慮以下幾種處理方式:

  • 提升 maxMessageSize 參數(shù)的大小,這樣可以容納更大的消息列表。但是,需要注意在提升參數(shù)大小時,要考慮到 RocketMQ broker 的性能和網(wǎng)絡(luò)帶寬等因素。

  • 考慮將消息列表進(jìn)行拆分,然后分批發(fā)送。這樣可以避免一次發(fā)送過多的消息。

  • 計算消息的大小并進(jìn)行壓縮??梢允褂靡恍嚎s算法,如 LZ4、Snappy 等,對消息進(jìn)行壓縮,以減小消息的大小。

  • 對超過 maxMessageSize 的消息進(jìn)行過濾或其他處理??梢酝ㄟ^業(yè)務(wù)邏輯對消息進(jìn)行分組或分類,對超過 maxMessageSize 的消息進(jìn)行過濾或其他處理,以避免發(fā)送超出限制的消息。

開發(fā)者在開發(fā)時需要注意消息列表的大小限制,避免出現(xiàn)超出限制的情況。

2 消息的批量消費

2.1 批量消費的優(yōu)點

批量消費消息可以提高 RocketMQ 的消費者性能和吞吐量,因為批量消費消息可以減少網(wǎng)絡(luò) I/O 操作和降低消息消費延遲。批量消費消息在以下情況下特別有用:

  • 消費大量小型消息時
  • 需要降低消息消費延遲時
  • 需要提高消費者性能時
    但是,批量消費消息也存在一些注意事項,需要注意以下幾點:
  • 消息列表的大小不能超過 broker 設(shè)置的最大消息大小。
  • 消息列表的順序可能與單條消息的順序不同。如果要保持消息順序,需要對消息進(jìn)行分組。
  • 批量消費消息會增大消息重試的難度。因此,如果的消息消費邏輯具有事務(wù)性質(zhì),建議使用單條消息消費方式。

批量消費消息是一種提高 RocketMQ 消費者性能和吞吐量的好方法,但需要注意消息列表大小、消息順序和事務(wù)性質(zhì)等問題,以確保消費者的可靠性和穩(wěn)定性。

2.2 推、拉和長輪詢

MQ的消費模式可以大致分為兩種,一種是推Push,一種是拉Pull

  • Push是服務(wù)端主動推送消息給客戶端,優(yōu)點是及時性較好,但如果客戶端沒有做好流控,一旦服務(wù)端推送大量消息到客戶端時,就會導(dǎo)致客戶端消息堆積甚至崩潰。
  • Pull是客戶端需要主動到服務(wù)端取數(shù)據(jù),優(yōu)點是客戶端可以依據(jù)自己的消費能力進(jìn)行消費,但拉取的頻率也需要用戶自己控制,拉取頻繁容易造成服務(wù)端和客戶端的壓力,拉取間隔長又容易造成消費不及時。

2.3 對pull(拉模式)的批量消費

DefaultLitePullConsumer是RocketMQ中的拉模式消息消費者,其工作流程如下:

  • 初始化:創(chuàng)建DefaultLitePullConsumer實例,并設(shè)置相關(guān)參數(shù),如消費者組、NameServer地址等。
    訂閱主題:調(diào)用DefaultLitePullConsumer的subscribe()方法,訂閱感興趣的主題(Topic)及消息過濾標(biāo)簽(Tag)。
  • 消費者啟動:調(diào)用DefaultLitePullConsumer的start()方法,啟動消費者。啟動后,消費者將與NameServer建立連接, 獲取感興趣的主題的路由信息,找到對應(yīng)的Broker,之后建立與Broker的連接。然后,消費者向Broker報告當(dāng)前的消費進(jìn)度。
  • 輪詢拉取消息:調(diào)用DefaultLitePullConsumer的poll()方法,輪詢地從Broker拉取消息。拉取消息后,本地需要有一個處理消息的邏輯。
  • 消息處理:處理拉取到的消息,如處理業(yè)務(wù)邏輯、持久化等。消息處理需要考慮潛在的并發(fā)和處理速度問題。
  • 提交消費進(jìn)度:當(dāng)消息處理完成后,調(diào)用DefaultLitePullConsumer的commitSync()或commitAsync()方法,將消費進(jìn)度提交給Broker,以便下次拉取時能繼續(xù)從上一次完成處理的消息開始拉取。
  • 異常處理:如果拉取消息過程中出現(xiàn)異常,可以考慮重試。簡單的重試可以通過DefaultLitePullConsumer提供的seek()方法回滾消費進(jìn)度,但需要注意處理消息時的冪等性問題。
  • 關(guān)閉消費者:當(dāng)需要關(guān)閉消費者時,調(diào)用DefaultLitePullConsumer的shutdown()方法。關(guān)閉過程中會斷開與Broker和NameServer的連接,并釋放相關(guān)資源。

在使用DefaultLitePullConsumer時,需要注意控制拉取消息的速率(比如使用定時任務(wù)調(diào)用poll()方法)及消息處理的并發(fā)能力。根據(jù)實際業(yè)務(wù)去實現(xiàn)適當(dāng)?shù)奶幚聿呗?,保證在消費者速率和處理能力之間達(dá)到一個平衡。

demo:

 public static void main(String[] args) throws Exception {
        DefaultLitePullConsumer litePullConsumer = new DefaultLitePullConsumer("lite_pull_consumer_test");
        litePullConsumer.setNamesrvAddr("localhost:9876");
        litePullConsumer.subscribe("test_topic", "*");
        litePullConsumer.setPullBatchSize(20);
        litePullConsumer.start();
        try {
            while (running) {
                List<MessageExt> messageExts = litePullConsumer.poll();
                System.out.printf("%s%n", messageExts);
            }
        } finally {
            litePullConsumer.shutdown();
        }
    }

首先還是初始化DefaultLitePullConsumer并設(shè)置ConsumerGroupName,調(diào)用subscribe方法訂閱topic并啟動。與Push Consumer不同的是,LitePullConsumer拉取消息調(diào)用的是輪詢poll接口,如果能拉取到消息則返回對應(yīng)的消息列表,否則返回null。通過setPullBatchSize可以設(shè)置每一次拉取的最大消息數(shù)量,此外如果不額外設(shè)置,LitePullConsumer默認(rèn)是自動提交位點。在subscribe模式下,同一個消費組下的多個LitePullConsumer會負(fù)載均衡消費。

2.4 對push(推模式)的批量消費

DefaultMQPushConsumer的工作流程如下:

  • 客戶端初始化:創(chuàng)建DefaultMQPushConsumer實例,并設(shè)置相關(guān)配置參數(shù),如消費者組、NameServer地址等。
  • 消費者訂閱:消費者通過調(diào)用DefaultMQPushConsumer的subscribe()方法訂閱感興趣的主題(Topic)以及消息的過濾標(biāo)簽(Tag)。訂閱的Topic信息最終將傳遞給Broker。
  • 注冊MessageListener:實現(xiàn)一個自定義的MessageListener作為消費者處理消息的監(jiān)聽器。這個監(jiān)聽器會在消費者收到新消息時被調(diào)用。
  • 消費者啟動:消費者調(diào)用DefaultMQPushConsumer的start()方法啟動。啟動時,消費者將與NameServer建立連接,獲取感興趣的主題的路由信息,找到對應(yīng)的Broker,之建立連接。隨后,消費者會定期向Broker發(fā)起心跳請求以維持連接。
  • Broker消息推送:當(dāng)Broker上有新的生產(chǎn)者消息,它會分配給消費者組內(nèi)的消費者進(jìn)行消費。Broker會將消息推送給消費者,消費者在收到新消息后會執(zhí)行相應(yīng)的MessageListener邏輯處理消息。
  • 消息消費確認(rèn):消費者成功處理消息后,將發(fā)送ACK確認(rèn)信息回Broker。Broker收到ACK確認(rèn)后,會將消息標(biāo)記為消費成功并標(biāo)記為已處理。如果消費失敗,可以選擇重試,根據(jù)配置設(shè)定的策略最終可能會進(jìn)入死信隊列。
  • 逐步消費:消費者會持續(xù)消費新消息,直到關(guān)閉消費者。
  • 消費者關(guān)閉:調(diào)用DefaultMQPushConsumer的shutdown()方法關(guān)閉消費者。關(guān)閉時,消費者將斷開與Broker的連接,并釋放相關(guān)資源。

請注意,在使用DefaultMQPushConsumer時,消費者的并發(fā)能力由MessageListener實現(xiàn)來保證。因此,在設(shè)計MessageListener實現(xiàn)時需要考慮到高并發(fā)處理能力。雖然RocketMQ客戶端提供了設(shè)置消費線程池的配置選項,但還是推薦根據(jù)實際需求來實現(xiàn)合適的并發(fā)方案。‘

2.41 使用DefaultMQPushConsumer :

批量消費消息需要在消費者端設(shè)置 ConsumeMessageBatchMaxSize 參數(shù),以指定每次批量消費的消息數(shù)量。

public static void main(String[] args) throws InterruptedException, MQClientException {
        // step1: 創(chuàng)建一個DefaultMQPushConsumer實例
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("BatchPushConsumer");
        consumer.setNamesrvAddr("localhost:9876");
        consumer.setConsumeMessageBatchMaxSize(10); // 設(shè)置每次批量消費的消息數(shù)量

        // step2: 為消費者訂閱一個Topic
        consumer.subscribe("test_topic", "*");

        // step3: 注冊一個MessageListenerConcurrently,并實現(xiàn)批量消費邏輯
        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(
                    List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
                System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msgs);
                // 處理批量消息
                for (MessageExt msg : msgs) {
                    // 在此處理每條消息,例如保存到數(shù)據(jù)庫等
                    System.out.println("msg : " + new String(msg.getBody()));
                }
                // 返回消費狀態(tài)
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });

        // step4: 啟動消費者
        consumer.start();
        System.out.println("Consumer started!");

        // 讓主線程等待以保持進(jìn)程不退出
        TimeUnit.SECONDS.sleep(60);
    }

在上述示例中,每次批量消費的消息數(shù)量被設(shè)置為10。每次推送過來的消息數(shù)量可能并不總是達(dá)到這個數(shù)字,但是它不會超過這個數(shù)量。如果希望調(diào)整批量大小,可以通過consumer.setConsumeMessageBatchMaxSize();修改這個值。

2.4.2 :springboot 通過@RocketMQMessageListener 完成消息消費:

:默認(rèn)的 RocketMQ Spring Boot Starter 并不支持直接設(shè)置批量消費模式,消息是一個個處理的。

對于RocketMQ Listener可以管理多個線程同時處理消息。可以有多個消息同時處理。通過將順序消息設(shè)置為并發(fā)模式并設(shè)置消費線程數(shù)。

import com.example.demo.MessageProcessor;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;

@Service
@RocketMQMessageListener(topic = "my-topic", consumerGroup = "myConsumerGroup", consumeMode = ConsumeMode.CONCURRENTLY, consumeThreadMax = 3)
public class MyConsumer implements RocketMQListener<MessageExt> {
    @Autowired
    private MessageProcessor messageProcessor;

    @Override
    public void onMessage(MessageExt messageExt) {
        messageProcessor.process(messageExt);
    }
}

在這個例子中,consumeMode = ConsumeMode.CONCURRENTLY表示消費者將啟用并發(fā)模式,而consumeThreadMax = 3將同時處理三個消息。注意,這不是真正意義上的批量消費,而是通過多線程來同時處理多個消息。要實現(xiàn)批量消費,需要進(jìn)一步處理MessageExt消息,這取決于的實際需求。例如,可以緩沖消息,等待足夠多的消息可用后一次性處理它們。

在設(shè)置 consumeThreadMax 參數(shù)時,請確保它不要過大,以避免系統(tǒng)資源過載。同時,優(yōu)化MessageProcessor中的相關(guān)邏輯,盡量減少處理每個消息所需的時間。通過這種機制,雖然無法實現(xiàn)真正意義上的批量消費,但仍然可以幫助提高消息處理的效率。

總結(jié)

以上為個人經(jīng)驗,希望能給大家一個參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • IntelliJ?idea報junit?no?tasks?available問題的解決辦法

    IntelliJ?idea報junit?no?tasks?available問題的解決辦法

    這篇文章主要給大家介紹了關(guān)于IntelliJ?idea報junit?no?tasks?available問題的解決辦法,文中通過圖文介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2023-11-11
  • SpringBoot配置多數(shù)據(jù)源的四種方式分享

    SpringBoot配置多數(shù)據(jù)源的四種方式分享

    在日常開發(fā)中我們都是以單個數(shù)據(jù)庫進(jìn)行開發(fā),在小型項目中是完全能夠滿足需求的,但是,當(dāng)我們牽扯到大型項目的時候,單個數(shù)據(jù)庫就難以承受用戶的CRUD操作,那么此時,我們就需要使用多個數(shù)據(jù)源進(jìn)行讀寫分離的操作,本文就給大家介紹SpringBoot配置多數(shù)據(jù)源的方式
    2023-07-07
  • 解決接口調(diào)用報錯newSocketStream(..)failed:Too?many?open?files問題

    解決接口調(diào)用報錯newSocketStream(..)failed:Too?many?open?files問題

    這篇文章主要介紹了解決接口調(diào)用報錯newSocketStream(..)failed:Too?many?open?files問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-07-07
  • SpringCloud的JPA連接PostgreSql的教程

    SpringCloud的JPA連接PostgreSql的教程

    這篇文章主要介紹了SpringCloud的JPA接入PostgreSql 教程,本文給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2021-06-06
  • java 進(jìn)制轉(zhuǎn)換實例詳解

    java 進(jìn)制轉(zhuǎn)換實例詳解

    這篇文章主要介紹了java 進(jìn)制轉(zhuǎn)換實例詳解的相關(guān)資料,需要的朋友可以參考下
    2017-04-04
  • java中構(gòu)造方法及this關(guān)鍵字的用法實例詳解(超詳細(xì))

    java中構(gòu)造方法及this關(guān)鍵字的用法實例詳解(超詳細(xì))

    大家都知道,java作為一門內(nèi)容豐富的編程語言,其中涉及的范圍是十分廣闊的,下面這篇文章主要給大家介紹了關(guān)于java中構(gòu)造方法及this關(guān)鍵字用法的相關(guān)資料,需要的朋友可以參考下
    2022-04-04
  • SpringBoot集成Activiti7工作流引擎的示例代碼

    SpringBoot集成Activiti7工作流引擎的示例代碼

    本文主要介紹了SpringBoot集成Activiti7工作流引擎的示例代碼,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2024-11-11
  • Java FutureTask解析與實戰(zhàn)指南

    Java FutureTask解析與實戰(zhàn)指南

    本文解析FutureTask,支持異步執(zhí)行、取消、結(jié)果獲取,結(jié)合ExecutorService實現(xiàn)高效任務(wù)管理,適用于并發(fā)編程中的響應(yīng)優(yōu)化與進(jìn)度監(jiān)控,感興趣的朋友一起看看吧
    2025-07-07
  • Java多線程中常見的幾個問題

    Java多線程中常見的幾個問題

    這篇文章主要介紹了Java多線程中常見的幾個問題 ,需要的朋友可以參考下
    2015-05-05
  • SpringBoot 模糊映射(Ambiguous mapping)報錯解決指南

    SpringBoot 模糊映射(Ambiguous mapping)報錯解決指南

    本文針對 Spring Boot 啟動時 “requestMappingHandlerMapping 模糊映射” 引發(fā)的 BeanCreationException,下面就來詳細(xì)的介紹一下該錯誤的解決,感興趣的可以了解一下
    2025-09-09

最新評論

泸定县| 东方市| 安宁市| 和平县| 寿光市| 林西县| 屯昌县| 五华县| 新宾| 故城县| 恩平市| 镇巴县| 巨鹿县| 闵行区| 缙云县| 许昌市| 潢川县| 都昌县| 周宁县| 东海县| 长岛县| 调兵山市| 自治县| 扎囊县| 保山市| 马山县| 津市市| 呼玛县| 调兵山市| 手机| 砚山县| 岳阳市| 保定市| 揭阳市| 阜南县| 盐池县| 凤山县| 东丰县| 巴马| 巧家县| 芮城县|