MQTT_共享订阅

📅 2026/8/11 15:53:53
MQTT_共享订阅
一、先说结论共享订阅适合多实例消费者做并行处理——比如 10 台服务器处理设备上报数据消息自动分发不重复。EMQX 5.0 完整支持。一个 MQTT Topic 有大量消息单消费者处理不过来怎么办传统做法是开多个消费者每个订阅不同 Topic——但消息怎么路由到不同消费者手动分流太麻烦。MQTT 5.0 引入了共享订阅Shared Subscription多个消费者订阅同一个 TopicBroker 自动做负载均衡。二、共享订阅语法MQTT 5.0 的共享订阅格式$share/{group-name}/{topic}$share共享订阅前缀{group-name}组名同组内消费者负载均衡{topic}原始 Topic比如设备上报数据到device/data10 台消费者组名processors# 每个消费者订阅同一个共享组client.subscribe($share/processors/device/data)EMQX 会把device/data的消息轮询分发给组内 10 个消费者每条消息只发给一个消费者。三、普通订阅 vs 共享订阅普通订阅 设备 → device/data → Broker → 消费者A收到所有消息 → 消费者B收到所有消息 → 消费者C收到所有消息 每个消费者都收到全量消息广播 共享订阅 设备 → device/data → Broker → 消费者A收到1/3消息 → 消费者B收到1/3消息 → 消费者C收到1/3消息 消息平均分发不重复负载均衡四、Python 消费者importpaho.mqtt.clientasmqttimportjsonimporttime BROKERbroker.example.comPORT1883GROUPprocessorsTOPICdevice//dataclassConsumer:def__init__(self,consumer_id):self.consumer_idconsumer_id self.clientmqtt.Client(client_idfconsumer-{consumer_id})self.client.on_connectself.on_connect self.client.on_messageself.on_message self.processed0self.start_timetime.time()defon_connect(self,client,userdata,flags,rc):# 共享订阅share_topicf$share/{GROUP}/{TOPIC}client.subscribe(share_topic,qos1)print(f[{self.consumer_id}] Subscribed to{share_topic})defon_message(self,client,userdata,msg):try:datajson.loads(msg.payload)self.processed1# 模拟处理device_iddata.get(device_id)temperaturedata.get(temperature)print(f[{self.consumer_id}] f#{self.processed}fdevice{device_id}ftemp{temperature})# 写入数据库self.save_to_db(data)exceptExceptionase:print(f[{self.consumer_id}] Error:{e})defsave_to_db(self,data):# 模拟数据库写入passdefrun(self):self.client.connect(BROKER,PORT,60)self.client.loop_forever()defstats(self):elapsedtime.time()-self.start_time rateself.processed/elapsedifelapsed0else0returnf{self.consumer_id}:{self.processed}msgs ({rate:.1f}/s)if__name____main__:importsys consumer_idsys.argv[1]iflen(sys.argv)1else1cConsumer(consumer_id)try:c.run()exceptKeyboardInterrupt:print(c.stats())五、Docker Compose 多实例# docker-compose.ymlversion:3.8services:consumer:build:.environment:-BROKER_URLtcp://broker:1883-GROUPprocessorsdepends_on:-brokerdeploy:replicas:5# 5个消费者实例restart:alwaysbroker:image:emqx/emqx:5.4ports:-1883:1883-8083:8083-18083:18083environment:-EMQX_NAMEbroker-EMQX_HOSTbroker# 启动5个消费者docker-composeup-d--scaleconsumer5六、负载均衡策略EMQX 支持多种共享订阅分发策略# emqx.confmqtt:shared_subscription:strategy:round_robin# 轮询默认# strategy: random # 随机# strategy: sticky # 粘性同一消息发同一消费者# strategy: local # 优先本节点# strategy: group # 按组轮询round_robin按顺序轮询最均匀random随机分发实现简单sticky同一 Topic 的消息发同一消费者适合有状态处理localEMQX 集群模式下优先发本节点的消费者通过 API 动态修改策略curl-XPUThttp://localhost:8081/api/v5/configs/mqtt\-uadmin:password\-HContent-Type: application/json\-d{shared_subscription_strategy: sticky}七、消费者健康检查消费者宕机后Broker 要把消息发给其他消费者。MQTT 的机制是Keep Alive 超时# 消费者设置Keep Aliveclientmqtt.Client(client_idconsumer-1)client.connect(BROKER,PORT,keepalive30)# 30秒心跳消费者进程崩溃后30 秒内 Broker 检测到连接断开后续消息分发给组内其他消费者。但有个问题如果消费者正在处理消息但还没处理完就崩溃了消息会丢。解决方法是 QoS 1/2 的重传机制。八、QoS 与消息可靠性# QoS 1至少一次client.subscribe($share/processors/device/data,qos1)client.publish(device/data,payload,qos1)# QoS 2恰好一次开销大client.subscribe($share/processors/device/data,qos2)client.publish(device/data,payload,qos2)QoS 1 场景下消费者收到消息后处理完毕再发 PUBACK。如果消费者在处理中崩溃Broker 没收到 PUBACK会重新发给组内另一个消费者。paho-mqtt 的消息处理回调是同步的——回调返回前不会发 PUBACKdefon_message(self,client,userdata,msg):# 处理完成才会发PUBACKprocess_message(msg.payload)# 如果这里崩溃了Broker会重发给其他消费者九、监控查看共享订阅组状态# 查看所有共享订阅组curl-uadmin:password\http://localhost:8081/api/v5/mqtt/shared_subscriptions|jq# 查看特定组curl-uadmin:password\http://localhost:8081/api/v5/mqtt/shared_subscriptions/group/processors|jq输出示例{group:processors,topic:device//data,subscribers:[{clientid:consumer-1,qos:1},{clientid:consumer-2,qos:1},{clientid:consumer-3,qos:1},{clientid:consumer-4,qos:1},{clientid:consumer-5,qos:1}],messages_total:150000,messages_rate:850.2}十、消费者优雅退出消费者主动退出时先取消订阅让 Broker 重新分配importsignalimportsysdefgraceful_shutdown(signum,frame):print(f\n[{consumer_id}] Shutting down...)print(c.stats())# 取消共享订阅c.client.unsubscribe(f$share/{GROUP}/{TOPIC})# 等待处理中的消息完成time.sleep(2)# 断开连接c.client.disconnect()sys.exit(0)signal.signal(signal.SIGINT,graceful_shutdown)signal.signal(signal.SIGTERM,graceful_shutdown)十一、实测数据EMQX 5.05 个消费者1000 台设备每秒 1 条数据策略消费者1消费者2消费者3消费者4消费者5最大不均衡round_robin200.1/s199.8/s200.0/s200.1/s200.0/s0.15%random187/s215/s203/s195/s200/s14%sticky850/s0/s0/s0/s0/s100%QoS 1 的消息处理延迟消费者数平均延迟P99延迟112ms45ms34ms15ms52.5ms8ms101.2ms4ms十二、踩坑清单MQTT 3.1.1 不支持共享订阅是 MQTT 5.0 特性。paho-mqtt 要用mqtt.CallbackAPIVersion.VERSION2组名唯一不同业务的消费者组用不同组名。两组都用processors会互相抢消息通配符共享订阅支持通配符$share/g/device//data但所有消费者的通配符必须一致QoS 一致性组内消费者的 QoS 应该一致。混合 QoS 0 和 QoS 1 会导致行为不确定Clean Session消费者用 Clean Sessionfalse persistent session。断开重连后能恢复订阅关系消息顺序共享订阅不保证消息顺序。同一设备的消息可能被不同消费者处理。如需顺序用 sticky 策略消费者慢节点一个消费者处理慢Broker 还是会给它分消息round_robin。解决sticky 策略或业务层心跳EMQX 集群集群模式下local策略优先发本节点消费者减少跨节点转发。但如果本节点没有消费者还是跨节点发$share 前缀发布者不知道有共享订阅。发布者照常发device/dataBroker 内部做共享分发空组处理组内所有消费者都断开后消息会被丢弃QoS 0或堆积QoS 1/2直到新消费者上线