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

本地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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 19:37:01