Python脚本通过MQTT向C#应用传输序列化消息时存在丢包问题求助
MQTT消息丢失问题排查与修复
嘿,我来帮你搞定这个Python发消息、C#收不全的问题!先拆解核心原因,再给你改好代码~
核心问题:Python端重复创建MQTT客户端,未确保消息发送完成
你当前的Python代码里,每次调用mqtt()函数都会新建一个MQTT客户端对象,连接Broker后就调用publish,但paho的publish是异步操作——函数执行完后,client对象会被Python的垃圾回收机制销毁,很多消息还没真正发送到Broker就被中断了,这就是为什么会丢6、11、16、17这些消息。
另外,你的C#代码有个小冗余:每次收到消息都新建一个队列,只存当前消息再遍历,完全没必要,直接处理收到的字节数组就行。
修复方案
1. 修改Python代码:复用MQTT客户端,确保消息发送完成
把客户端的创建和连接逻辑放到函数外面,复用同一个客户端,同时启动网络循环处理异步操作,还可以等待publish确认消息已送达Broker:
import paho.mqtt.client as paho import AIML_pb2 import datetime as dt import time def on_publish(client, userdata, result): print(f"消息已发布,返回结果: {result}\n") # 全局复用MQTT客户端,避免重复创建连接 client = paho.Client("control1") client.on_publish = on_publish def publish_mqtt(buf): broker = "127.0.0.1" port = 1883 # 检查客户端是否已连接,未连接则尝试重连 if not client.is_connected(): try: client.connect(broker, port) # 启动MQTT网络循环,处理异步消息收发 client.loop_start() except Exception as e: print(f"连接Broker失败: {e}") return try: # 设置QoS=1,确保消息至少被Broker接收一次 publish_result = client.publish( "/SERVICES/RESPONSE/ECG/UDCMGR/AIML", buf, qos=1 ) # 等待消息发送完成,确保Broker确认接收 publish_result.wait_for_publish() if publish_result.rc != 0: print(f"消息发布失败,错误码: {publish_result.rc}") except Exception as e: print(f"发布消息时出错: {e}") def main(): detections = AIML_pb2.ObjectDetectionParametersBroadcast() # 提前初始化连接 publish_mqtt(None) for i in range(0, 20): detections.timestamp = dt.datetime.now().strftime('%Y-%m-%d_%H-%M-%S') detections.frameID = i detections.objectType = "cell phone" detections.conf = 0.90 detections.top_x = 0.5 detections.top_y = 0.6 detections.bottom_x = 0.7 detections.bottom_y = 0.9 buf = detections.SerializeToString() publish_mqtt(buf) print(f"已发送 frameID: {i}") time.sleep(1) # 程序结束后停止网络循环并断开连接 client.loop_stop() client.disconnect() if __name__ == "__main__": main()
2. 优化C#代码:去掉冗余队列,正确反序列化Protobuf消息
你原来的C#代码没有真正反序列化Protobuf消息,只是打印了字节数组的类型,现在改成正确解析消息并打印frameID:
if (e.Topic == "/SERVICES/RESPONSE/ECG/UDCMGR/AIML") { try { // 反序列化Protobuf消息 var detections = AIML_pb2.ObjectDetectionParametersBroadcast.Parser.ParseFrom(e.Message); Console.WriteLine($"{detections.FrameID} 成功解析Protobuf消息,类型: {e.Message.GetType()}"); ScriptAPIs.aimlStatus = true; } catch (Exception ex) { Console.WriteLine($"解析Protobuf消息失败: {ex.Message}"); } }
额外建议
- 设置QoS等级:Python和C#两端都使用QoS=1,确保消息至少被送达一次(默认QoS=0是“最多一次”,容易丢消息)。
- 检查Broker日志:查看你的MQTT Broker(比如Mosquitto)的日志,确认Broker是否收到了所有20条消息,进一步定位问题是发送端还是接收端。
- C#订阅时保持连接:确保C#客户端没有意外断开连接,订阅后持续监听消息。
内容的提问来源于stack exchange,提问作者elgatovolador88
相关产品推荐
相关产品推荐

