使用QoS1及以上发布时订阅者仅接收2/3消息的问题排查
MQTT QoS 1/2模式下末尾消息丢失问题排查
我用两个Python程序测试部署在Linux Docker中的Mosquitto Broker,Python程序运行于Windows机器:
- 一个订阅
Test主题,持续监听并将接收的所有消息打印到标准输出,同时写入日志文件 - 另一个尽可能快地向
Test主题发送指定数量的消息
测试现象
- QoS 0时,发布者和订阅者均正常,无消息丢失
- 发布者使用QoS 1时,所有消息显示已发送,但发送50条及以上消息时,订阅者仅能收到60%-90%的消息;测试100、1000、10000条消息场景时,无论发送数量多少,仅丢失末尾消息,中间消息无丢失,收到的消息约为总数的2/3
- 发布者改用QoS 2后情况更糟,订阅者仅能收到30%-50%的消息
- 订阅者QoS设为0/1/2均无变化,两个脚本以同一用户身份连接
疑惑:猜测系统可能无法及时处理所有消息,但不理解为何仅丢失末尾消息而非中间消息?
发布者代码
import paho.mqtt.client as mqtt def on_connect(client, userdata, flags, reason_code, properties): print(f"Connected with result code {reason_code}") broker_hostname = "myhost" port = 1883 client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2) client.username_pw_set(username="user", password="password") client.on_connect = on_connect client.connect(broker_hostname, port) client.loop_start() topic = "Test" msg_count = 0 try: while msg_count < 1000: msg_count += 1 result = client.publish(topic, msg_count, qos=1) status = result[0] if status == 0: print("Message " + str(msg_count) + " is published to topic " + topic) else: print("Failed to send message to topic " + topic) if not client.is_connected(): print("Client not connected, exiting...") break finally: client.disconnect() client.loop_stop()
订阅者代码
import paho.mqtt.client as mqtt import logging logging.basicConfig(filename="log.txt", level=logging.INFO) def on_connect(client, userdata, flags, reason_code, properties): print(f"Connected with result code {reason_code}") client.subscribe("Test", qos=0) # also tested with qos=2, no difference in behaviour def on_message(client, userdata, msg): print(msg.topic + " " + str(msg.payload)) logging.info(msg.payload) broker_hostname = "myhost" port = 1883 client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2) client.username_pw_set(username="user", password="password") client.on_connect = on_connect client.on_message = on_message client.connect(broker_hostname, port) client.loop_forever()
Mosquitto配置
persistence true persistence_location /mosquitto/data/ log_type subscribe log_type unsubscribe log_type websockets log_type error log_type warning log_type notice log_type information log_dest file /mosquitto/log/mosquitto.log log_dest stdout password_file /mosquitto/passwd_file allow_anonymous false # MQTT Default listener listener 1883 0.0.0.0 # MQTT over WebSockets listener 9001 0.0.0.0 protocol websockets
问题原因分析
核心问题出在发布者的关闭逻辑:你直接调用client.disconnect()断开连接,但QoS 1/2的消息需要Broker返回确认(PUBACK/PUBCOMP),发布者的loop_start()线程在loop_stop()后会立即终止,导致还未收到确认的末尾消息被直接丢弃,Broker也不会重试这些消息(因为客户端已断开)。
只丢末尾消息的原因:前面的消息在快速发送过程中已经收到了Broker的确认,而末尾的消息刚发送出去,还没来得及完成确认流程,客户端就被强制断开了。
另外,publish返回的状态0仅表示消息成功加入客户端的发送队列,不代表Broker已经确认接收,这也是你看到“所有消息都已发送”但实际丢失的原因。
解决方法
1. 等待所有未完成的QoS消息确认后再断开
修改发布者的finally块,先等待所有消息的确认完成,再停止循环并断开连接:
import time # 记得添加该导入 # ... 原有代码 ... finally: # 等待所有QoS 1/2消息的确认完成 while client.in_publish() > 0: time.sleep(0.1) client.disconnect() client.loop_stop()
2. 调整发布速率,避免消息堆积
如果发送速度过快,客户端和Broker的确认流程跟不上,也会导致未确认消息在断开时丢失。可以在每次publish后添加短暂延迟:
result = client.publish(topic, msg_count, qos=1) time.sleep(0.01) # 给确认流程留处理时间
3. 订阅者QoS匹配(可选)
虽然你测试过订阅者QoS无变化,但确保订阅者QoS不低于发布者QoS,能保证Broker按更高的QoS级别处理消息传递,避免不必要的降级。比如发布者用QoS 1,订阅者也设为QoS 1。
内容的提问来源于stack exchange,提问作者Jan
相关产品推荐
相关产品推荐

