RocketMQ?消息Message的結(jié)構(gòu)和使用方式詳解
?? RocketMQ 消息(Message)詳解
在 Apache RocketMQ 中,消息(Message) 是數(shù)據(jù)傳輸?shù)淖钚卧?,是生產(chǎn)者與消費(fèi)者之間通信的“載體”。理解 Message 的結(jié)構(gòu)、屬性、生命周期和使用方式,是掌握 RocketMQ 的核心基礎(chǔ)。
推薦閱讀:深入理解Apache RocketMQ 中Message 消息的核心概念
一、什么是 Message?
? 定義:
Message 是 RocketMQ 中封裝實(shí)際業(yè)務(wù)數(shù)據(jù)的對(duì)象,包含消息體(Body)和一系列元數(shù)據(jù)(如 Topic、Tag、Key、Properties 等),用于在生產(chǎn)者與消費(fèi)者之間傳遞信息。
類(lèi)比:就像一封信,信紙是內(nèi)容(Body),信封上寫(xiě)著收件人(Topic)、標(biāo)簽(Tag)、編號(hào)(Key)等信息。
二、Message 的核心結(jié)構(gòu)
一個(gè) Message 對(duì)象主要由以下幾個(gè)部分組成:
| 字段 | 類(lèi)型 | 是否必填 | 說(shuō)明 |
|---|---|---|---|
| Topic | String | ? 必填 | 消息所屬的主題,用于路由和分類(lèi) |
| Body | byte[] | ? 必填 | 消息的實(shí)際內(nèi)容,通常為序列化后的 JSON、Protobuf 等 |
| Tags | String | ? 可選 | 子分類(lèi)標(biāo)簽,用于消費(fèi)者過(guò)濾(如 CREATE, CANCEL) |
| Keys | String | ? 可選 | 消息的唯一鍵或業(yè)務(wù)主鍵(如訂單號(hào)),用于排查、索引 |
| Flag | int | ? 可選 | 消息標(biāo)志位(如是否壓縮) |
| DelayTimeLevel | int | ? 可選 | 延遲消息級(jí)別(1~18),實(shí)現(xiàn)定時(shí)投遞 |
| Properties | Map<String, String> | ? 可選 | 自定義屬性,RocketMQ 內(nèi)部也使用它存儲(chǔ)系統(tǒng)屬性 |
三、Message 各字段詳解
1.Topic(主題)
- 消息的邏輯分類(lèi),決定消息被發(fā)送到哪個(gè)隊(duì)列。
- 必須提前創(chuàng)建或允許自動(dòng)創(chuàng)建。
- 示例:
ORDER_TOPIC,USER_LOG_TOPIC
new Message("ORDER_TOPIC", ...);
2.Body(消息體)
- 實(shí)際傳輸?shù)臄?shù)據(jù),必須是字節(jié)數(shù)組。
- 通常通過(guò) JSON、Protobuf、Hessian 等序列化框架編碼。
String content = "{\"orderId\":\"1001\",\"userId\":10086}";
Message msg = new Message(topic, tag, content.getBytes(StandardCharsets.UTF_8));
?? 注意:
- 單條消息大小默認(rèn)最大 4MB(可配置)
- 過(guò)大消息會(huì)影響性能,建議拆分或使用外部存儲(chǔ)(如上傳文件后傳 URL)
3.Tags(標(biāo)簽)
- 用于對(duì)同一 Topic 下的消息進(jìn)行二次分類(lèi)。
- 消費(fèi)者可通過(guò)
subscribe("Topic", "TagA || TagB")進(jìn)行過(guò)濾。
// 發(fā)送
new Message("ORDER_TOPIC", "CREATE", "創(chuàng)建訂單".getBytes());
new Message("ORDER_TOPIC", "PAY", "支付完成".getBytes());
// 訂閱 CREATE 類(lèi)型消息
consumer.subscribe("ORDER_TOPIC", "CREATE");? 優(yōu)勢(shì):輕量級(jí)過(guò)濾,避免消費(fèi)者接收無(wú)關(guān)消息。
?? 注意:Tags 是字符串匹配,不支持正則(但支持
*通配和||多選)
4.Keys(消息鍵)
- 為消息設(shè)置唯一標(biāo)識(shí)或業(yè)務(wù)主鍵(如訂單號(hào)、用戶ID)。
- 支持通過(guò)
mqadmin queryMsgByKey命令查詢(xún)消息。 - 支持索引,便于問(wèn)題排查。
Message msg = new Message(...);
msg.setKeys("ORDER_20240501001");
? 建議:關(guān)鍵業(yè)務(wù)消息務(wù)必設(shè)置 Keys,便于追蹤。
5.Properties(屬性)
- 鍵值對(duì)形式的擴(kuò)展字段,可用于:
- 存儲(chǔ)自定義上下文(如 traceId、tenantId)
- RocketMQ 內(nèi)部使用(如
RECONSUME_TIME、DELAY、TRAN_MSG)
msg.putUserProperty("traceId", "abc123");
msg.putUserProperty("source", "web");
?? 注意:系統(tǒng)屬性以
PREFIX_SYS_PROP開(kāi)頭,不要沖突。
6.DelayTimeLevel(延遲級(jí)別)
- 設(shè)置消息延遲投遞時(shí)間,實(shí)現(xiàn)“定時(shí)任務(wù)”功能。
- 取值范圍:1~18,對(duì)應(yīng)不同延遲時(shí)間:
| 級(jí)別 | 時(shí)間 |
|---|---|
| 1 | 1s |
| 2 | 5s |
| 3 | 10s |
| 4 | 30s |
| 5 | 1m |
| 6 | 2m |
| 7 | 3m |
| 8 | 4m |
| 9 | 5m |
| 10 | 6m |
| 11 | 7m |
| 12 | 8m |
| 13 | 9m |
| 14 | 10m |
| 15 | 20m |
| 16 | 30m |
| 17 | 1h |
| 18 | 2h |
Message msg = new Message("DELAY_TOPIC", "TAG", "延遲消息".getBytes());
msg.setDelayTimeLevel(5); // 延遲1分鐘
producer.send(msg);
?? 注意:延遲消息不保證精確時(shí)間,存在輕微誤差。
四、Message 的生命周期
1. 生產(chǎn)者創(chuàng)建 Message 對(duì)象 ↓ 2. 發(fā)送到 Broker(寫(xiě)入 CommitLog) ↓ 3. 構(gòu)建 ConsumeQueue 和 IndexFile ↓ 4. 消費(fèi)者拉取消息(根據(jù) Topic + Queue) ↓ 5. 處理成功 → 提交 Offset ↓ 6. 消息過(guò)期(默認(rèn) 72 小時(shí))→ 被刪除
? 消息是持久化存儲(chǔ)的,即使消費(fèi)者未上線,消息也不會(huì)丟失。
五、Message 的存儲(chǔ)機(jī)制
雖然 Message 是邏輯對(duì)象,但在 Broker 端有嚴(yán)格的物理存儲(chǔ)結(jié)構(gòu):
1.CommitLog
- 所有消息按到達(dá)順序追加寫(xiě)入 CommitLog 文件(順序?qū)?,高性能?/li>
- 每個(gè)消息包含:Topic、Queue、Body、Properties 等完整信息
2.ConsumeQueue
- 每個(gè) Topic 的每個(gè) MessageQueue 對(duì)應(yīng)一個(gè) ConsumeQueue
- 存儲(chǔ)消息的邏輯偏移量、大小、物理位置,用于快速定位消息
ConsumeQueue/{Topic}/{QueueId}/
├── 00000000000000000000
└── ...
3.IndexFile
- 可選索引文件,支持通過(guò) Keys 或時(shí)間范圍 查詢(xún)消息
- 用于排查問(wèn)題(如“查找某個(gè)訂單的消息”)
IndexFile/index_1714567890000
六、Message 的發(fā)送方式回顧
| 方式 | 說(shuō)明 |
|---|---|
| 同步發(fā)送 | 阻塞等待結(jié)果,適用于關(guān)鍵消息 |
| 異步發(fā)送 | 回調(diào)通知結(jié)果,高吞吐場(chǎng)景 |
| 單向發(fā)送 | 不關(guān)心結(jié)果,日志類(lèi)消息 |
| 事務(wù)消息 | 半消息 + 本地事務(wù) + 提交/回滾 |
所有方式發(fā)送的都是
Message對(duì)象。
七、最佳實(shí)踐與注意事項(xiàng)
| 實(shí)踐 | 說(shuō)明 |
|---|---|
| ? 設(shè)置 Topic 和 Tag | 合理分類(lèi),便于管理和過(guò)濾 |
| ? 關(guān)鍵消息設(shè)置 Keys | 便于通過(guò) mqadmin 查詢(xún) |
| ? 控制 Body 大小 | ≤ 4MB,避免影響性能 |
| ? 使用 UTF-8 編碼 | 防止亂碼 |
| ? 避免空 Body | 可能導(dǎo)致異常 |
| ? 合理使用延遲消息 | 替代部分定時(shí)任務(wù),但不要濫用 |
? 自定義屬性用 putUserProperty | 避免覆蓋系統(tǒng)屬性 |
八、常見(jiàn)問(wèn)題排查
| 問(wèn)題 | 原因 | 解決方案 |
|---|---|---|
MessageExt is null | 拉取超時(shí)或無(wú)消息 | 正?,F(xiàn)象,重試即可 |
msg put message to store error | 消息過(guò)大或磁盤(pán)滿 | 檢查大小限制和磁盤(pán)空間 |
| 延遲消息未按時(shí)投遞 | 時(shí)間誤差或 Broker 壓力大 | 接受輕微延遲,或使用外部調(diào)度系統(tǒng) |
| 通過(guò) Key 查不到消息 | IndexFile 未生成或已過(guò)期 | 檢查 messageIndexEnable 配置 |
| 消息重復(fù) | 網(wǎng)絡(luò)重試、Rebalance | 消費(fèi)者做冪等處理 |
? 總結(jié):Message 核心要點(diǎn)
| 維度 | 說(shuō)明 |
|---|---|
| 角色 | 消息傳輸?shù)幕締卧?/td> |
| 組成 | Topic + Body + Tags + Keys + Properties + Delay |
| 大小限制 | 默認(rèn) ≤ 4MB |
| 存儲(chǔ)方式 | 順序?qū)?CommitLog,索引通過(guò) ConsumeQueue 和 IndexFile |
| 可查詢(xún)性 | 支持按 Key、時(shí)間、Offset 查詢(xún) |
| 擴(kuò)展性 | 支持自定義屬性,靈活傳遞上下文 |
| 高級(jí)功能 | 支持延遲、事務(wù)、順序消息 |
?? 一句話總結(jié):
Message 是 RocketMQ 的“數(shù)據(jù)包” —— 它不僅是業(yè)務(wù)數(shù)據(jù)的載體,更是路由、過(guò)濾、追蹤、延遲、事務(wù)等功能的基礎(chǔ)。
設(shè)計(jì)好 Message 的結(jié)構(gòu)與屬性,才能讓消息系統(tǒng)真正高效、可靠、易維護(hù)。
掌握 Message,你就掌握了 RocketMQ 的“語(yǔ)言”。
到此這篇關(guān)于RocketMQ 消息Message的結(jié)構(gòu)和使用方式詳解的文章就介紹到這了,更多相關(guān)RocketMQ 消息Message內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Spring MVC學(xué)習(xí)之DispatcherServlet請(qǐng)求處理詳析
這篇文章主要給大家介紹了關(guān)于Spring MVC學(xué)習(xí)教程之DispatcherServlet請(qǐng)求處理的相關(guān)資料,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面來(lái)一起學(xué)習(xí)學(xué)習(xí)吧2018-11-11
如何使用Java?8函數(shù)式編程優(yōu)雅處理多層嵌套數(shù)據(jù)
Java8是Java語(yǔ)言歷史上的一個(gè)重大更新,它帶來(lái)了許多新的特性和改進(jìn),其中函數(shù)式編程的引入是其亮點(diǎn)之一,這篇文章主要介紹了如何使用Java?8函數(shù)式編程優(yōu)雅處理多層嵌套數(shù)據(jù)的相關(guān)資料,需要的朋友可以參考下2026-01-01
一文搞懂Java MD5算法的原理及實(shí)現(xiàn)
MD5信息摘要算法,一種被廣泛使用的密碼散列函數(shù),可以產(chǎn)生出一個(gè)128位(16字節(jié))的散列值(hash value),用于確保信息傳輸完整一致。本文將詳解MD5算法的原理及實(shí)現(xiàn),感興趣的可以了解一下2022-06-06
MyBatisPlus 主鍵策略的實(shí)現(xiàn)(4種)
MyBatis Plus 集成了多種主鍵策略,幫助用戶快速生成主鍵,本文主要介紹了MyBatisPlus主鍵策略的實(shí)現(xiàn),具有一定的參考價(jià)值,感興趣的可以了解一下2023-10-10

