本地Docker部署Mosquitto时QoS=1订阅仅收20条消息求助
MQTT订阅端仅接收20条记录的排查与解决
核心原因及解决方案
1. Mosquitto Broker的Inflight消息限制
Mosquitto默认限制每个客户端的未确认QoS 1/2消息数量为20条。当发布端每秒发送500条QoS=1消息时,若订阅端未及时返回PUBACK确认,Broker会暂停推送新消息,直到收到之前消息的确认。
解决:
- 修改Mosquitto配置文件
mosquitto.conf,添加/调整以下参数:max_inflight_messages 1000 # 允许更多未确认消息 max_queued_messages 5000 # 增大消息队列容量 - 重启Broker容器:
docker restart <你的mosquitto容器名>
2. 订阅端消息处理阻塞
如果订阅端的消息回调函数包含同步阻塞操作(如慢IO、复杂计算),会导致无法及时发送PUBACK,Broker因此停止推送新消息。
解决:
- 将消息处理逻辑放到异步线程/进程中,让回调函数快速返回。以Python paho-mqtt为例的优化代码:
import paho.mqtt.client as mqtt import threading def on_message(client, userdata, msg): # 异步处理消息,避免阻塞回调 threading.Thread(target=handle_msg, args=(msg.payload,)).start() def handle_msg(payload): # 这里编写消息解析、存储等业务逻辑 print(f"Received: {payload.decode()}") client = mqtt.Client() client.on_message = on_message client.connect("localhost", 1883, 60) client.subscribe("test/topic", qos=1) client.loop_forever()
3. 发布端发送速率失控
若发布端不控制发送节奏,Broker的消息堆积可能触发限流机制,导致订阅端无法接收全量消息。
解决:
- 发布端加入流量控制,等待部分
PUBACK后再继续发送,避免瞬间压垮Broker和订阅端:import paho.mqtt.client as mqtt import time import threading ack_count = 0 lock = threading.Lock() def on_publish(client, userdata, mid): global ack_count with lock: ack_count += 1 client = mqtt.Client() client.on_publish = on_publish client.connect("localhost", 1883, 60) client.loop_start() total_messages = 500 * 20 sent_count = 0 while sent_count < total_messages: # 控制未确认消息不超过100条 with lock: if sent_count - ack_count > 100: time.sleep(0.001) continue client.publish("test/topic", f"msg_{sent_count}", qos=1) sent_count += 1 # 维持每秒500条的发送速率 if sent_count % 500 == 0: time.sleep(1) client.loop_stop() client.disconnect()
4. Docker容器资源不足
如果Mosquitto容器分配的CPU、内存资源不足,会导致消息转发延迟,进而影响订阅端接收效率。
解决:
- 查看容器资源占用情况:
docker stats <你的mosquitto容器名> - 重启容器时分配更多资源:
docker run -d --name mosquitto --cpus 2 --memory 512M eclipse-mosquitto
内容的提问来源于stack exchange,提问作者Ritab
相关产品推荐
相关产品推荐

