列機(jī)制:售貨柜多設(shè)備消息隔離設(shè)計(jì))
主題Topic與隊(duì)列機(jī)制售貨柜多設(shè)備消息隔離設(shè)計(jì)作者黒漂技術(shù)佬系列專欄RocketMQ核心原理與無人售貨柜項(xiàng)目實(shí)戰(zhàn)一、Topic消息的一級(jí)分類1.1 Topic是什么Topic主題是RocketMQ中最頂層的消息分類單位。你可以把它類比為數(shù)據(jù)庫里的表——消息是表里的行每條消息都屬于某張表。數(shù)據(jù)庫類比 數(shù)據(jù)庫 → RocketMQ 數(shù)據(jù)庫 → Broker 表 → Topic 行 → Message 分區(qū) → Queue 列標(biāo)簽 → Tag一個(gè)Broker上可以存放多個(gè)Topic每個(gè)Topic存放一類業(yè)務(wù)消息。比如售貨柜項(xiàng)目里Topic名用途生產(chǎn)者消費(fèi)者order_topic訂單消息訂單服務(wù)庫存服務(wù)、推送服務(wù)payment_topic支付消息支付服務(wù)柜子網(wǎng)關(guān)、ERP同步服務(wù)device_topic設(shè)備消息柜子網(wǎng)關(guān)監(jiān)控服務(wù)、告警服務(wù)log_topic日志消息各服務(wù)日志收集服務(wù)1.2 Topic的創(chuàng)建方式自動(dòng)創(chuàng)建Producer第一次向一個(gè)不存在的Topic發(fā)消息時(shí)Broker會(huì)自動(dòng)創(chuàng)建它。開發(fā)環(huán)境方便但生產(chǎn)環(huán)境強(qiáng)烈建議關(guān)閉autoCreateTopicEnablefalse原因有三自動(dòng)創(chuàng)建的Topic默認(rèn)4個(gè)隊(duì)列可能不符合業(yè)務(wù)需求容易因拼寫錯(cuò)誤創(chuàng)建出錯(cuò)誤的Topicorder_topicvsorder_topc消息發(fā)到錯(cuò)誤地方排查困難自動(dòng)創(chuàng)建的Topic均勻分布在所有Broker上不可控手動(dòng)創(chuàng)建通過Dashboard或命令行預(yù)先創(chuàng)建可以指定隊(duì)列數(shù)、所在Broker等# 命令行創(chuàng)建Topicshmqadmin updateTopic\-n127.0.0.1:9876\-b127.0.0.1:10911\-torder_topic\-r8\# 讀隊(duì)列數(shù)-w8# 寫隊(duì)列數(shù)也可以用Dashboard界面操作主題 → 新增 → 填寫Topic名和隊(duì)列數(shù)。1.3 讀寫隊(duì)列的含義創(chuàng)建Topic時(shí)會(huì)指定讀隊(duì)列數(shù)r和寫隊(duì)列數(shù)w這倆有什么區(qū)別寫隊(duì)列WriteQueueProducer發(fā)消息時(shí)Broker按寫隊(duì)列數(shù)做路由分配消息實(shí)際存在這些隊(duì)列里讀隊(duì)列ReadQueueConsumer消費(fèi)時(shí)按讀隊(duì)列數(shù)做負(fù)載均衡從這些隊(duì)列拉消息正常情況下讀隊(duì)列數(shù) 寫隊(duì)列數(shù)。什么時(shí)候會(huì)不一樣Topic縮容。假設(shè)原來8個(gè)隊(duì)列想縮到4個(gè)直接改寫隊(duì)列數(shù)為4讀隊(duì)列數(shù)暫時(shí)保持8等Consumer把原8個(gè)隊(duì)列的消息消費(fèi)完再把讀隊(duì)列數(shù)改成4。這樣縮容不會(huì)丟消息。二、QueueTopic下的子分區(qū)2.1 Queue的作用Queue隊(duì)列是Topic的子分區(qū)類似數(shù)據(jù)庫表的分區(qū)。一個(gè)Topic默認(rèn)有4個(gè)隊(duì)列可配置。Topic: order_topic (4個(gè)隊(duì)列) Queue-0 ──→ [msg1] [msg5] [msg9] ... Queue-1 ──→ [msg2] [msg6] [msg10] ... Queue-2 ──→ [msg3] [msg7] [msg11] ... Queue-3 ──→ [msg4] [msg8] [msg12] ...Producer發(fā)消息時(shí)默認(rèn)輪詢Round Robin把消息均勻分配到各隊(duì)列。Queue的兩個(gè)核心作用并行消費(fèi)多個(gè)Consumer可以分別消費(fèi)不同Queue實(shí)現(xiàn)并行處理。1個(gè)Topic有8個(gè)Queue最多8個(gè)Consumer同時(shí)消費(fèi)吞吐量線性擴(kuò)展。負(fù)載均衡ConsumerGroup內(nèi)的Consumer實(shí)例自動(dòng)分配Queue誰消費(fèi)哪個(gè)Queue由Rebalance算法決定。某個(gè)Consumer掛了它的Queue會(huì)被重新分配給其他Consumer。2.2 隊(duì)列數(shù)怎么定隊(duì)列數(shù)不是越多越好也不是越少越好。經(jīng)驗(yàn)法則場(chǎng)景建議隊(duì)列數(shù)原因低頻消息訂單4~8消費(fèi)者實(shí)例少多了也用不上中頻消息設(shè)備狀態(tài)8~16多個(gè)區(qū)域消費(fèi)者并行高頻消息日志/埋點(diǎn)16~32高并發(fā)需要更多并行度順序消息按業(yè)務(wù)分區(qū)鍵數(shù)量定保證同一Key的消息在同一Queue售貨柜項(xiàng)目建議訂單Topic 8個(gè)隊(duì)列8個(gè)消費(fèi)實(shí)例夠用設(shè)備消息Topic 16個(gè)隊(duì)列按區(qū)域分配日志Topic 32個(gè)隊(duì)列高吞吐。三、Tag消息的二級(jí)分類3.1 Tag的概念Tag標(biāo)簽是Topic下的二級(jí)分類用于在同一個(gè)Topic內(nèi)區(qū)分子類消息。如果Topic是數(shù)據(jù)庫的表那Tag就是表里的一個(gè)分類字段。Topic: device_topic ├── Tag: heartbeat 設(shè)備心跳消息 ├── Tag: alert 設(shè)備告警消息 ├── Tag: status 設(shè)備狀態(tài)消息 └── Tag: inventory 設(shè)備庫存消息為什么不用多個(gè)Topic代替Tag因?yàn)門opic是物理隔離每個(gè)Topic占獨(dú)立的存儲(chǔ)和隊(duì)列資源。用Tag在同一Topic下分類共享隊(duì)列資源減少Topic數(shù)量降低管理成本。3.2 Tag的使用Producer端指定Tag// Topic:Tag 格式rocketMQTemplate.syncSend(device_topic:heartbeat,heartbeatMsg);rocketMQTemplate.syncSend(device_topic:alert,alertMsg);rocketMQTemplate.syncSend(device_topic:status,statusMsg);Consumer端按Tag過濾消費(fèi)// 只消費(fèi)告警消息RocketMQMessageListener(topicdevice_topic,selectorExpressionalert,// 只消費(fèi)Tagalert的消息consumerGroupalert_consumer_group)publicclassAlertConsumerimplementsRocketMQListenerAlertMessage{OverridepublicvoidonMessage(AlertMessagemessage){alertService.handle(message);}}// 消費(fèi)心跳和狀態(tài)消息多Tag用 || 分隔RocketMQMessageListener(topicdevice_topic,selectorExpressionheartbeat || status,consumerGroupmonitor_consumer_group)publicclassMonitorConsumerimplementsRocketMQListenerMessageExt{OverridepublicvoidonMessage(MessageExtmessage){Stringtagmessage.getTags();if(heartbeat.equals(tag)){handleHeartbeat(message);}elseif(status.equals(tag)){handleStatus(message);}}}// 消費(fèi)所有TagRocketMQMessageListener(topicdevice_topic,selectorExpression*,// *表示消費(fèi)所有TagconsumerGroupall_device_consumer_group)publicclassAllDeviceConsumerimplementsRocketMQListenerMessageExt{// ...}3.3 Tag vs Topic的選擇標(biāo)準(zhǔn)什么時(shí)候用不同Topic什么時(shí)候用不同Tag記住一個(gè)原則消費(fèi)方不同、需要物理隔離→ 用不同Topic消費(fèi)方相同或部分相同、邏輯分類→ 用同一Topic 不同Tag舉例場(chǎng)景選擇原因訂單消息 vs 支付消息不同Topic消費(fèi)方完全不同物理隔離設(shè)備心跳 vs 設(shè)備告警同Topic不同Tag都屬于設(shè)備消息監(jiān)控服務(wù)都要消費(fèi)支付成功 vs 支付失敗同Topic不同Tag都是支付消息下游消費(fèi)邏輯接近四、ConsumerGroup與隊(duì)列分配關(guān)系4.1 隊(duì)列分配規(guī)則在集群消費(fèi)模式下一個(gè)ConsumerGroup內(nèi)的多個(gè)Consumer實(shí)例分?jǐn)俆opic的所有Queue。核心規(guī)則一個(gè)Queue同一時(shí)間只被組內(nèi)一個(gè)Consumer實(shí)例消費(fèi)。Topic: order_topic (4個(gè)Queue) ConsumerGroup: order_consumer_group 情況12個(gè)Consumer實(shí)例 Consumer-1 ← Queue-0, Queue-1 Consumer-2 ← Queue-2, Queue-3 情況24個(gè)Consumer實(shí)例 Consumer-1 ← Queue-0 Consumer-2 ← Queue-1 Consumer-3 ← Queue-2 Consumer-4 ← Queue-3 情況36個(gè)Consumer實(shí)例超過隊(duì)列數(shù) Consumer-1 ← Queue-0 Consumer-2 ← Queue-1 Consumer-3 ← Queue-2 Consumer-4 ← Queue-3 Consumer-5 ← 空閑分不到隊(duì)列 Consumer-6 ← 空閑分不到隊(duì)列4.2 消費(fèi)者超過隊(duì)列數(shù)怎么辦如上所示當(dāng)Consumer實(shí)例數(shù) Queue數(shù)時(shí)多出來的Consumer空閑不消費(fèi)任何消息。這不是Bug是設(shè)計(jì)如此——Queue是并行消費(fèi)的最小單位4個(gè)Queue最多4個(gè)Consumer并行。所以部署消費(fèi)服務(wù)時(shí)實(shí)例數(shù)不要超過Topic的Queue數(shù)否則浪費(fèi)資源。如果需要更多并行度先增加Queue數(shù)。4.3 Rebalance機(jī)制ConsumerGroup內(nèi)的Consumer實(shí)例數(shù)變化時(shí)擴(kuò)容/縮容/宕機(jī)RocketMQ會(huì)自動(dòng)觸發(fā)Rebalance重平衡重新分配Queue。初始狀態(tài) Consumer-1 ← Queue-0, Queue-1 Consumer-2 ← Queue-2, Queue-3 Consumer-2宕機(jī) → 觸發(fā)Rebalance Consumer-1 ← Queue-0, Queue-1, Queue-2, Queue-3 全部接管 新Consumer-3加入 → 觸發(fā)Rebalance Consumer-1 ← Queue-0, Queue-1 Consumer-3 ← Queue-2, Queue-3Rebalance由Consumer端發(fā)起每20秒檢查一次。如果發(fā)現(xiàn)隊(duì)列分配發(fā)生變化自動(dòng)調(diào)整。這個(gè)過程對(duì)用戶透明但有一個(gè)注意點(diǎn)Rebalance瞬間可能出現(xiàn)消息重復(fù)投遞Consumer切換隊(duì)列時(shí)上一次未確認(rèn)的消息會(huì)被重新投遞所以消費(fèi)端一定要做冪等。五、售貨柜多設(shè)備消息隔離實(shí)戰(zhàn)方案5.1 問題背景假設(shè)我們有以下業(yè)務(wù)需求全國(guó)有10000臺(tái)售貨柜分布在500個(gè)門店每臺(tái)柜子定時(shí)上報(bào)心跳、庫存、狀態(tài)柜子關(guān)門后上報(bào)訂單消息柜子異常時(shí)上報(bào)告警消息不同門店的消息需要隔離處理A店的運(yùn)維只關(guān)心A店的設(shè)備柜子出貨消息要保證同一臺(tái)設(shè)備的順序性5.2 隔離方案設(shè)計(jì)方案一按門店ID區(qū)分TopicTopic: store_10001_device_topic (門店10001的設(shè)備消息) Topic: store_10002_device_topic (門店10002的設(shè)備消息) ...優(yōu)點(diǎn)物理隔離徹底不同門店互不影響缺點(diǎn)500個(gè)門店 500個(gè)TopicTopic數(shù)量爆炸管理成本高RocketMQ建議單Broker Topic數(shù)不超過5000但太多影響性能方案二按設(shè)備ID分配隊(duì)列 消息Key這是推薦的方案。用統(tǒng)一的Topic通過Queue分配和消息Key來實(shí)現(xiàn)邏輯隔離Topic: device_message (16個(gè)Queue) ├── 用Tag區(qū)分消息類型heartbeat / alert / status / inventory ├── 用設(shè)備ID作為消息Key便于查詢 └── 用MessageQueueSelector把同一設(shè)備的消息路由到同一QueueProducer端路由ServicepublicclassDeviceMessageService{ResourceprivateRocketMQTemplaterocketMQTemplate;/** * 發(fā)送設(shè)備消息同一設(shè)備的消息路由到同一隊(duì)列保證順序 */publicvoidsendDeviceMessage(StringdeviceId,Stringtag,Objectpayload){DeviceMessagemessagenewDeviceMessage(deviceId,tag,payload);// 使用hashKey路由同一deviceId的消息始終進(jìn)入同一QueuerocketMQTemplate.syncSendOrderly(device_message:tag,// Topic:TagMessageBuilder.withPayload(message).build(),deviceId// hashKey按設(shè)備ID做hash選隊(duì)列);}}syncSendOrderly方法內(nèi)部用MessageQueueSelector對(duì) deviceId 取hash后對(duì)隊(duì)列數(shù)取模保證同一設(shè)備的消息始終進(jìn)同一隊(duì)列。這樣同一設(shè)備的消息被同一Consumer消費(fèi)保證了消息順序性。Consumer端按門店過濾ComponentRocketMQMessageListener(topicdevice_message,selectorExpressionalert || status,// 只消費(fèi)告警和狀態(tài)consumerGroupstore_monitor_group,consumeModeConsumeMode.CONCURRENTLY)publicclassStoreMonitorConsumerimplementsRocketMQListenerDeviceMessage{OverridepublicvoidonMessage(DeviceMessagemessage){StringdeviceIdmessage.getDeviceId();// 從設(shè)備ID查出所屬門店StringstoreIddeviceService.getStoreId(deviceId);// 按門店分發(fā)處理StoreHandlerhandlerstoreHandlerMap.get(storeId);if(handler!null){handler.handle(message);}}}方案三按消息類型用Tag區(qū)分 按區(qū)域用ConsumerGroupTopic: device_message Tag: heartbeat → ConsumerGroup: heartbeat_group (全國(guó)心跳匯總) Tag: alert → ConsumerGroup: alert_group_north (北方區(qū)域告警) ConsumerGroup: alert_group_south (南方區(qū)域告警) Tag: status → ConsumerGroup: status_group (狀態(tài)監(jiān)控) Tag: inventory → ConsumerGroup: inventory_group (庫存同步)不同ConsumerGroup各自消費(fèi)全量消息在Consumer內(nèi)部按區(qū)域/門店過濾處理。這種方式靈活但ConsumerGroup多注意不要超過RocketMQ的訂閱組限制默認(rèn)1000個(gè)。5.3 最終推薦方案綜合考慮售貨柜項(xiàng)目的消息隔離方案如下┌─────────────────────────────────────────────────────────┐ │ Topic 設(shè)計(jì) │ ├─────────────────────────────────────────────────────────┤ │ │ │ order_topic (8隊(duì)列) │ │ └─ Tag: order_created / order_paid / order_closed │ │ │ │ device_message (16隊(duì)列) │ │ └─ Tag: heartbeat / alert / status / inventory │ │ └─ hashKey: deviceId (保證同設(shè)備消息順序) │ │ │ │ payment_callback (8隊(duì)列) │ │ └─ Tag: wechat / alipay │ │ │ │ device_log (32隊(duì)列) │ │ └─ Tag: operation / error / access │ │ └─ 單向發(fā)送不走順序 │ │ │ ├─────────────────────────────────────────────────────────┤ │ ConsumerGroup 設(shè)計(jì) │ ├─────────────────────────────────────────────────────────┤ │ │ │ order_topic: │ │ inventory_consumer_group (庫存服務(wù)集群模式) │ │ push_consumer_group (推送服務(wù)集群模式) │ │ │ │ device_message: │ │ alert_consumer_group (告警服務(wù)) │ │ monitor_consumer_group (監(jiān)控服務(wù)消費(fèi)heartbeatstatus) │ │ inventory_sync_group (庫存同步服務(wù)消費(fèi)inventory) │ │ │ │ payment_callback: │ │ gateway_consumer_group (柜子網(wǎng)關(guān)消費(fèi)后通知出貨) │ │ erp_sync_consumer_group (ERP同步服務(wù)) │ │ │ └─────────────────────────────────────────────────────────┘5.4 關(guān)鍵設(shè)計(jì)決策總結(jié)設(shè)計(jì)決策選擇理由門店隔離方式消息Key Consumer內(nèi)過濾避免Topic爆炸邏輯隔離夠用設(shè)備消息順序hashKeydeviceId路由到同一Queue出貨和庫存變動(dòng)需保序消息類型區(qū)分Tag同類設(shè)備消息共享Topic減少Topic數(shù)消費(fèi)并行度Queue數(shù) 預(yù)計(jì)最大Consumer實(shí)例數(shù)避免實(shí)例空閑浪費(fèi)冪等保障訂單ID/設(shè)備ID時(shí)間戳做去重防止Rebalance導(dǎo)致重復(fù)消費(fèi)日志類消息獨(dú)立Topic 單向發(fā)送和業(yè)務(wù)消息隔離互不影響六、小結(jié)這一篇從Topic、Queue、Tag三個(gè)維度拆解了RocketMQ的消息分類和分區(qū)機(jī)制重點(diǎn)講解了Queue的并行消費(fèi)和負(fù)載均衡作用、Tag的二級(jí)分類過濾、ConsumerGroup與Queue的分配關(guān)系。最后給出了一套完整的售貨柜多設(shè)備消息隔離方案按業(yè)務(wù)域分Topic、按消息類型分Tag、按設(shè)備ID做Queue路由保證順序、按消費(fèi)方分ConsumerGroup。這套方案在后面的系列文章中會(huì)持續(xù)用到。