You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 00:17:16