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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 07:52:35