Debian11 Bullseye下Paho MQTT大消息传输异常及配置咨询
MQTT图片传输异常排查与配置指南
问题背景
- 项目基于MQTT协议实现图片传输,采用mosquitto作为消息broker,发送端运行Debian GNU/Linux 11 (bullseye)操作系统。
- 基于paho-mqtt 1.6.1开发的客户端基础功能正常,可正常发送字符串类型payload,但发送bytearray格式图片消息时出现异常:
publish()方法返回元组(0,1)提示消息发送成功,但broker端、订阅对应主题的其他客户端均未收到任何消息。 - 对照测试结果:在Raspberry Pi 4设备上运行完全相同的Python脚本、发送同一张测试图片时功能完全正常,两台设备paho-mqtt版本一致。
- 初步排查确认问题与消息大小阈值相关:Debian 11环境下发布长度超过1931字节的字符串时会触发相同异常,发送400字节以内的图片可正常传输。
测试用客户端代码
图片发送端代码
import paho.mqtt.client as mqtt, time BROKER = "***.***.*.***" CAMID = "Clinux" class MyMqtt: def __init__(self, broker, clientId): self.broker = broker self.clientId = clientId def run(self): self.client = mqtt.Client(client_id=self.clientId) self.client.loop_start() connectionSuccess = False while connectionSuccess == False: try: self.client.connect(self.broker) connectionSuccess = True print("Connected!\nbroker: ",self.broker) except: print("Unable to connect to ", self.broker) time.sleep(1) self.client.on_connect = self.on_connect self.client.on_disconnect = self.on_disconnect self.isConnected = True def Disconnect(self): self.client.disconnect() self.client.loop_stop() def on_connect(self, client, userdata, flags, rc): self.Image() def on_disconnect(self, client, userdata, rc): if rc != 0: print("Unexpected disconnection.") else: print("Successfully Disconnected") self.isConnected = False def Publish(self, topic, message, qos = 0, retain = False): return self.client.publish(topic, message, qos, retain) def Subscribe(self, topic, callbackFunction): self.client.subscribe(topic) self.client.message_callback_add(topic, callbackFunction) def Image(self): try: f=open("sample.png", "rb") #3.7kiB in same folder fileContent = f.read() byteArr = bytearray(fileContent) print(self.Publish("Image", byteArr)) except Exception as e: print("erro") if __name__ == '__main__': mqttClient = MyMqtt(BROKER, CAMID) mqttClient.run() while mqttClient.isConnected == True: print("on loop") time.sleep(1)
图片接收端代码
import paho.mqtt.client as mqtt, time BROKER = "***.***.*.***" CAMID = "CPC" class MyMqtt: def __init__(self, broker, clientId): self.broker = broker self.clientId = clientId def run(self): self.client = mqtt.Client(self.clientId) self.client.loop_start() connectionSuccess = False while connectionSuccess == False: try: self.client.connect(self.broker) connectionSuccess = True print("Connected!\nbroker: ",self.broker) except: print("Unable to connect to ", self.broker) time.sleep(1) self.client.on_connect = self.on_connect self.client.on_disconnect = self.on_disconnect self.isConnected = True def Disconnect(self): self.client.disconnect() self.client.loop_stop() def on_connect(self, client, userdata, flags, rc): self.Subscribe("Image", self.ReceiveImage) def on_disconnect(self, client, userdata, rc): if rc != 0: print("Unexpected disconnection.") else: print("Successfully Disconnected") self.isConnected = False def Publish(self, topic, message, qos = 0, retain = False): self.client.publish(topic, message, qos, retain) def Subscribe(self, topic, callbackFunction): self.client.subscribe(topic) self.client.message_callback_add(topic, callbackFunction) def ReceiveImage(self, client, userdata, message): try: f = open('output.jpg', "wb") f.write(message.payload) print("Image Received") f.close() except Exception as e: print(e.message) if __name__ == '__main__': mqttClient = MyMqtt(BROKER, CAMID) mqttClient.run() while mqttClient.isConnected == True: print("on loop") time.sleep(1)
消息大小限制配置方案
该异常和不同Linux发行版的二进制构建逻辑无关,是配置项不匹配、代码时序错误导致的,按以下步骤调整即可:
1. Mosquitto Broker端配置
打开mosquitto配置文件(默认路径/etc/mosquitto/mosquitto.conf),添加/修改消息大小限制项:
# 单条MQTT消息最大允许大小,单位为字节,示例值为10MB,可根据实际传输的图片大小调整 message_size_limit 10485760
保存后执行systemctl restart mosquitto重启broker服务生效。
2. 客户端代码逻辑修正
现有代码存在回调绑定时序问题:所有回调函数都在connect()调用成功之后才绑定,会导致连接建立后的报文流控、分片逻辑初始化异常,大消息发送时会出现假成功。
修正逻辑:在调用connect()方法之前完成所有回调函数绑定,修正后的run()方法示例:
def run(self): self.client = mqtt.Client(client_id=self.clientId) # 所有回调必须在connect前绑定 self.client.on_connect = self.on_connect self.client.on_disconnect = self.on_disconnect # 可选:调整客户端队列、飞行窗口大小,适配大消息传输 self.client.max_inflight_messages_set(100) self.client.max_queued_messages_set(1000) self.client.loop_start() connectionSuccess = False while connectionSuccess == False: try: self.client.connect(self.broker) connectionSuccess = True print("Connected!\nbroker: ",self.broker) except: print("Unable to connect to ", self.broker) time.sleep(1) self.isConnected = True
发送端、接收端代码都需要做上述时序调整。
3. 链路MTU/MSS排查
如果调整完上述配置仍存在1900字节左右的阈值,检查网卡、防火墙、路由设备的TCP分段限制,执行以下命令调整链路参数:
# 将eth0替换为实际使用的联网网卡名 ip link set dev eth0 mtu 1500 iptables -t mangle -A POSTROUTING -p tcp --tcp-flags SYN,RST SYN -j TCPMSS --clamp-mss-to-pmtu
内容的提问来源于stack exchange,提问作者Arthur Zim
相关产品推荐
相关产品推荐

