Rocketmq消息批量發(fā)送&消息批量消費方式
前言:批量發(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問題的解決辦法
這篇文章主要給大家介紹了關(guān)于IntelliJ?idea報junit?no?tasks?available問題的解決辦法,文中通過圖文介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考借鑒價值,需要的朋友可以參考下2023-11-11
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問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教2024-07-07
SpringCloud的JPA連接PostgreSql的教程
這篇文章主要介紹了SpringCloud的JPA接入PostgreSql 教程,本文給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下2021-06-06
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工作流引擎的示例代碼,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2024-11-11
SpringBoot 模糊映射(Ambiguous mapping)報錯解決指南
本文針對 Spring Boot 啟動時 “requestMappingHandlerMapping 模糊映射” 引發(fā)的 BeanCreationException,下面就來詳細(xì)的介紹一下該錯誤的解決,感興趣的可以了解一下2025-09-09

