
主题Topic与队列机制售货柜多设备消息隔离设计作者黒漂技术佬系列专栏RocketMQ核心原理与无人售货柜项目实战一、Topic消息的一级分类1.1 Topic是什么Topic主题是RocketMQ中最顶层的消息分类单位。你可以把它类比为数据库里的表——消息是表里的行每条消息都属于某张表。数据库类比 数据库 → RocketMQ 数据库 → Broker 表 → Topic 行 → Message 分区 → Queue 列标签 → Tag一个Broker上可以存放多个Topic每个Topic存放一类业务消息。比如售货柜项目里Topic名用途生产者消费者order_topic订单消息订单服务库存服务、推送服务payment_topic支付消息支付服务柜子网关、ERP同步服务device_topic设备消息柜子网关监控服务、告警服务log_topic日志消息各服务日志收集服务1.2 Topic的创建方式自动创建Producer第一次向一个不存在的Topic发消息时Broker会自动创建它。开发环境方便但生产环境强烈建议关闭autoCreateTopicEnablefalse原因有三自动创建的Topic默认4个队列可能不符合业务需求容易因拼写错误创建出错误的Topicorder_topicvsorder_topc消息发到错误地方排查困难自动创建的Topic均匀分布在所有Broker上不可控手动创建通过Dashboard或命令行预先创建可以指定队列数、所在Broker等# 命令行创建Topicshmqadmin updateTopic\-n127.0.0.1:9876\-b127.0.0.1:10911\-torder_topic\-r8\# 读队列数-w8# 写队列数也可以用Dashboard界面操作主题 → 新增 → 填写Topic名和队列数。1.3 读写队列的含义创建Topic时会指定读队列数r和写队列数w这俩有什么区别写队列WriteQueueProducer发消息时Broker按写队列数做路由分配消息实际存在这些队列里读队列ReadQueueConsumer消费时按读队列数做负载均衡从这些队列拉消息正常情况下读队列数 写队列数。什么时候会不一样Topic缩容。假设原来8个队列想缩到4个直接改写队列数为4读队列数暂时保持8等Consumer把原8个队列的消息消费完再把读队列数改成4。这样缩容不会丢消息。二、QueueTopic下的子分区2.1 Queue的作用Queue队列是Topic的子分区类似数据库表的分区。一个Topic默认有4个队列可配置。Topic: order_topic (4个队列) Queue-0 ──→ [msg1] [msg5] [msg9] ... Queue-1 ──→ [msg2] [msg6] [msg10] ... Queue-2 ──→ [msg3] [msg7] [msg11] ... Queue-3 ──→ [msg4] [msg8] [msg12] ...Producer发消息时默认轮询Round Robin把消息均匀分配到各队列。Queue的两个核心作用并行消费多个Consumer可以分别消费不同Queue实现并行处理。1个Topic有8个Queue最多8个Consumer同时消费吞吐量线性扩展。负载均衡ConsumerGroup内的Consumer实例自动分配Queue谁消费哪个Queue由Rebalance算法决定。某个Consumer挂了它的Queue会被重新分配给其他Consumer。2.2 队列数怎么定队列数不是越多越好也不是越少越好。经验法则场景建议队列数原因低频消息订单4~8消费者实例少多了也用不上中频消息设备状态8~16多个区域消费者并行高频消息日志/埋点16~32高并发需要更多并行度顺序消息按业务分区键数量定保证同一Key的消息在同一Queue售货柜项目建议订单Topic 8个队列8个消费实例够用设备消息Topic 16个队列按区域分配日志Topic 32个队列高吞吐。三、Tag消息的二级分类3.1 Tag的概念Tag标签是Topic下的二级分类用于在同一个Topic内区分子类消息。如果Topic是数据库的表那Tag就是表里的一个分类字段。Topic: device_topic ├── Tag: heartbeat 设备心跳消息 ├── Tag: alert 设备告警消息 ├── Tag: status 设备状态消息 └── Tag: inventory 设备库存消息为什么不用多个Topic代替Tag因为Topic是物理隔离每个Topic占独立的存储和队列资源。用Tag在同一Topic下分类共享队列资源减少Topic数量降低管理成本。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过滤消费// 只消费告警消息RocketMQMessageListener(topicdevice_topic,selectorExpressionalert,// 只消费Tagalert的消息consumerGroupalert_consumer_group)publicclassAlertConsumerimplementsRocketMQListenerAlertMessage{OverridepublicvoidonMessage(AlertMessagemessage){alertService.handle(message);}}// 消费心跳和状态消息多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);}}}// 消费所有TagRocketMQMessageListener(topicdevice_topic,selectorExpression*,// *表示消费所有TagconsumerGroupall_device_consumer_group)publicclassAllDeviceConsumerimplementsRocketMQListenerMessageExt{// ...}3.3 Tag vs Topic的选择标准什么时候用不同Topic什么时候用不同Tag记住一个原则消费方不同、需要物理隔离→ 用不同Topic消费方相同或部分相同、逻辑分类→ 用同一Topic 不同Tag举例场景选择原因订单消息 vs 支付消息不同Topic消费方完全不同物理隔离设备心跳 vs 设备告警同Topic不同Tag都属于设备消息监控服务都要消费支付成功 vs 支付失败同Topic不同Tag都是支付消息下游消费逻辑接近四、ConsumerGroup与队列分配关系4.1 队列分配规则在集群消费模式下一个ConsumerGroup内的多个Consumer实例分摊Topic的所有Queue。核心规则一个Queue同一时间只被组内一个Consumer实例消费。Topic: order_topic (4个Queue) ConsumerGroup: order_consumer_group 情况12个Consumer实例 Consumer-1 ← Queue-0, Queue-1 Consumer-2 ← Queue-2, Queue-3 情况24个Consumer实例 Consumer-1 ← Queue-0 Consumer-2 ← Queue-1 Consumer-3 ← Queue-2 Consumer-4 ← Queue-3 情况36个Consumer实例超过队列数 Consumer-1 ← Queue-0 Consumer-2 ← Queue-1 Consumer-3 ← Queue-2 Consumer-4 ← Queue-3 Consumer-5 ← 空闲分不到队列 Consumer-6 ← 空闲分不到队列4.2 消费者超过队列数怎么办如上所示当Consumer实例数 Queue数时多出来的Consumer空闲不消费任何消息。这不是Bug是设计如此——Queue是并行消费的最小单位4个Queue最多4个Consumer并行。所以部署消费服务时实例数不要超过Topic的Queue数否则浪费资源。如果需要更多并行度先增加Queue数。4.3 Rebalance机制ConsumerGroup内的Consumer实例数变化时扩容/缩容/宕机RocketMQ会自动触发Rebalance重平衡重新分配Queue。初始状态 Consumer-1 ← Queue-0, Queue-1 Consumer-2 ← Queue-2, Queue-3 Consumer-2宕机 → 触发Rebalance Consumer-1 ← Queue-0, Queue-1, Queue-2, Queue-3 全部接管 新Consumer-3加入 → 触发Rebalance Consumer-1 ← Queue-0, Queue-1 Consumer-3 ← Queue-2, Queue-3Rebalance由Consumer端发起每20秒检查一次。如果发现队列分配发生变化自动调整。这个过程对用户透明但有一个注意点Rebalance瞬间可能出现消息重复投递Consumer切换队列时上一次未确认的消息会被重新投递所以消费端一定要做幂等。五、售货柜多设备消息隔离实战方案5.1 问题背景假设我们有以下业务需求全国有10000台售货柜分布在500个门店每台柜子定时上报心跳、库存、状态柜子关门后上报订单消息柜子异常时上报告警消息不同门店的消息需要隔离处理A店的运维只关心A店的设备柜子出货消息要保证同一台设备的顺序性5.2 隔离方案设计方案一按门店ID区分TopicTopic: store_10001_device_topic (门店10001的设备消息) Topic: store_10002_device_topic (门店10002的设备消息) ...优点物理隔离彻底不同门店互不影响缺点500个门店 500个TopicTopic数量爆炸管理成本高RocketMQ建议单Broker Topic数不超过5000但太多影响性能方案二按设备ID分配队列 消息Key这是推荐的方案。用统一的Topic通过Queue分配和消息Key来实现逻辑隔离Topic: device_message (16个Queue) ├── 用Tag区分消息类型heartbeat / alert / status / inventory ├── 用设备ID作为消息Key便于查询 └── 用MessageQueueSelector把同一设备的消息路由到同一QueueProducer端路由ServicepublicclassDeviceMessageService{ResourceprivateRocketMQTemplaterocketMQTemplate;/** * 发送设备消息同一设备的消息路由到同一队列保证顺序 */publicvoidsendDeviceMessage(StringdeviceId,Stringtag,Objectpayload){DeviceMessagemessagenewDeviceMessage(deviceId,tag,payload);// 使用hashKey路由同一deviceId的消息始终进入同一QueuerocketMQTemplate.syncSendOrderly(device_message:tag,// Topic:TagMessageBuilder.withPayload(message).build(),deviceId// hashKey按设备ID做hash选队列);}}syncSendOrderly方法内部用MessageQueueSelector对 deviceId 取hash后对队列数取模保证同一设备的消息始终进同一队列。这样同一设备的消息被同一Consumer消费保证了消息顺序性。Consumer端按门店过滤ComponentRocketMQMessageListener(topicdevice_message,selectorExpressionalert || status,// 只消费告警和状态consumerGroupstore_monitor_group,consumeModeConsumeMode.CONCURRENTLY)publicclassStoreMonitorConsumerimplementsRocketMQListenerDeviceMessage{OverridepublicvoidonMessage(DeviceMessagemessage){StringdeviceIdmessage.getDeviceId();// 从设备ID查出所属门店StringstoreIddeviceService.getStoreId(deviceId);// 按门店分发处理StoreHandlerhandlerstoreHandlerMap.get(storeId);if(handler!null){handler.handle(message);}}}方案三按消息类型用Tag区分 按区域用ConsumerGroupTopic: device_message Tag: heartbeat → ConsumerGroup: heartbeat_group (全国心跳汇总) Tag: alert → ConsumerGroup: alert_group_north (北方区域告警) ConsumerGroup: alert_group_south (南方区域告警) Tag: status → ConsumerGroup: status_group (状态监控) Tag: inventory → ConsumerGroup: inventory_group (库存同步)不同ConsumerGroup各自消费全量消息在Consumer内部按区域/门店过滤处理。这种方式灵活但ConsumerGroup多注意不要超过RocketMQ的订阅组限制默认1000个。5.3 最终推荐方案综合考虑售货柜项目的消息隔离方案如下┌─────────────────────────────────────────────────────────┐ │ Topic 设计 │ ├─────────────────────────────────────────────────────────┤ │ │ │ order_topic (8队列) │ │ └─ Tag: order_created / order_paid / order_closed │ │ │ │ device_message (16队列) │ │ └─ Tag: heartbeat / alert / status / inventory │ │ └─ hashKey: deviceId (保证同设备消息顺序) │ │ │ │ payment_callback (8队列) │ │ └─ Tag: wechat / alipay │ │ │ │ device_log (32队列) │ │ └─ Tag: operation / error / access │ │ └─ 单向发送不走顺序 │ │ │ ├─────────────────────────────────────────────────────────┤ │ ConsumerGroup 设计 │ ├─────────────────────────────────────────────────────────┤ │ │ │ order_topic: │ │ inventory_consumer_group (库存服务集群模式) │ │ push_consumer_group (推送服务集群模式) │ │ │ │ device_message: │ │ alert_consumer_group (告警服务) │ │ monitor_consumer_group (监控服务消费heartbeatstatus) │ │ inventory_sync_group (库存同步服务消费inventory) │ │ │ │ payment_callback: │ │ gateway_consumer_group (柜子网关消费后通知出货) │ │ erp_sync_consumer_group (ERP同步服务) │ │ │ └─────────────────────────────────────────────────────────┘5.4 关键设计决策总结设计决策选择理由门店隔离方式消息Key Consumer内过滤避免Topic爆炸逻辑隔离够用设备消息顺序hashKeydeviceId路由到同一Queue出货和库存变动需保序消息类型区分Tag同类设备消息共享Topic减少Topic数消费并行度Queue数 预计最大Consumer实例数避免实例空闲浪费幂等保障订单ID/设备ID时间戳做去重防止Rebalance导致重复消费日志类消息独立Topic 单向发送和业务消息隔离互不影响六、小结这一篇从Topic、Queue、Tag三个维度拆解了RocketMQ的消息分类和分区机制重点讲解了Queue的并行消费和负载均衡作用、Tag的二级分类过滤、ConsumerGroup与Queue的分配关系。最后给出了一套完整的售货柜多设备消息隔离方案按业务域分Topic、按消息类型分Tag、按设备ID做Queue路由保证顺序、按消费方分ConsumerGroup。这套方案在后面的系列文章中会持续用到。