MQTT发布端发送数据但订阅端未接收(QoS2异常)问题排查
MQTT不同QoS等级下数据传输异常排查
环境配置
- Raspberry Pi 3B+
- Win10 x64
- Paho MQTT客户端
- Mosquitto MQTT代理
问题现象
通过MQTT从树莓派向Win10设备每秒发送大小为358.4 Kb的负载,不同QoS等级下表现各异:
- QoS 0:可接收大部分数据,但存在部分负载丢失,发送与接收计数不匹配;
- QoS 1:可接收大部分数据,但存在部分负载丢失,发送与接收计数不匹配;
- QoS 2:发布端显示已完成全部数据发送,但订阅端未接收到任何内容,两端实时计数均能体现该异常。
代码实现
发布端(C语言)
相关代码片段如下:
#include "../daqhats/examples/c/daqhats_utils.h" // 其他头文件省略 #include "MQTTClient.h" #define MQTT_ADDRSS "tcp://localhost:1883" #define CLIENTID "ExampleClientPub" #define TOPIC "mqtt_test" #define QOS 2 #define TIMEOUT 1000L int pub_count = 0; int publish(MQTTClient client, MQTTClient_deliveryToken token, MQTTClient_connectOptions conn_opts, char* payload, int device) { int rc = 0; if ((rc = MQTTClient_connect(client, &conn_opts)) != MQTTCLIENT_SUCCESS) { printf("Failed to connect, return code %d\n", rc); exit(-1); } int payload_len = 25600 * sizeof(double); MQTTClient_publish(client, TOPIC, payload_len, &payload, QOS, 0, &token); printf("Published: %d\n", device); pub_count++; return rc; } int main(void) { // MQTT初始化 MQTTClient client; MQTTClient_connectOptions conn_opts = MQTTClient_connectOptions_initializer; MQTTClient_deliveryToken token; MQTTClient_create(&client, MQTT_ADDRSS, CLIENTID, MQTTCLIENT_PERSISTENCE_NONE, NULL); conn_opts.keepAliveInterval = 20; conn_opts.cleansession = 1; union DAQ_Data{ double f[2 * 25600]; char s[2 * 25600 * sizeof(double)]; }; // DAQ硬件配置代码省略 union DAQ_Data payload; do //while (is_running) { for (device = 0; device < DEVICE_COUNT; device++) { // 读取数据 result = mcc172_a_in_scan_read(address[device], &scan_status[device], samples_to_read, timeout, payload.f, buffer_size, &samples_read[device]); publish(client, token, conn_opts, payload.s, device); // 其他DAQ操作代码省略 } fflush(stdout); usleep(100000); if (enter_press()) { printf("Aborted - enter pressed\n\n"); break; } } while (!enter_press()); printf("Pub count: %d\n", pub_count); sleep(2); // DAQ硬件清理代码省略 return 0; }
订阅端(Python)
基于Eclipse MQTT Client示例开发,代码如下:
import argparse import os import ssl import sys import paho.mqtt.client as mqtt parser = argparse.ArgumentParser() parser.add_argument('-H', '--host', required=False, default="192.168.30.212")#"mqtt.eclipseprojects.io") parser.add_argument('-t', '--topic', required=False, default="mqtt_test") #"$SYS/#") parser.add_argument('-q', '--qos', required=False, type=int, default=2) parser.add_argument('-c', '--clientid', required=False, default=None) parser.add_argument('-u', '--username', required=False, default=None) parser.add_argument('-d', '--disable-clean-session', action='store_true', help="disable 'clean session' (sub + msgs not cleared when client disconnects)") parser.add_argument('-p', '--password', required=False, default=None) parser.add_argument('-P', '--port', required=False, type=int, default=None, help='Defaults to 8883 for TLS or 1883 for non-TLS') parser.add_argument('-k', '--keepalive', required=False, type=int, default=60) parser.add_argument('-s', '--use-tls', action='store_true') parser.add_argument('--insecure', action='store_true') parser.add_argument('-F', '--cacerts', required=False, default=None) parser.add_argument('--tls-version', required=False, default=None, help='TLS protocol version, can be one of tlsv1.2 tlsv1.1 or tlsv1\n') parser.add_argument('-D', '--debug', action='store_true') args, unknown = parser.parse_known_args() mq = open("mq_test.bin", "wb") rec_count = 0 def on_connect(mqttc, obj, flags, rc): print("rc: " + str(rc)) def on_message(mqttc, obj, msg): global rec_count rec_count = rec_count + 1 # message = msg.payload.decode("utf-8") + '\n' mq.write(msg.payload) print("recieved: ", rec_count) def on_publish(mqttc, obj, mid): print("mid: " + str(mid)) def on_subscribe(mqttc, obj, mid, granted_qos): print("Subscribed: " + str(mid) + " " + str(granted_qos)) def on_log(mqttc, obj, level, string): print(string) usetls = args.use_tls if args.cacerts: usetls = True port = args.port if port is None: if usetls: port = 8883 else: port = 1883 mqttc = mqtt.Client(args.clientid,clean_session = not args.disable_clean_session) if usetls: if args.tls_version == "tlsv1.2": tlsVersion = ssl.PROTOCOL_TLSv1_2 elif args.tls_version == "tlsv1.1": tlsVersion = ssl.PROTOCOL_TLSv1_1 elif args.tls_version == "tlsv1": tlsVersion = ssl.PROTOCOL_TLSv1 elif args.tls_version is None: tlsVersion = None else: print ("Unknown TLS version - ignoring") tlsVersion = None if not args.insecure: cert_required = ssl.CERT_REQUIRED else: cert_required = ssl.CERT_NONE mqttc.tls_set(ca_certs=args.cacerts, certfile=None, keyfile=None, cert_reqs=cert_required, tls_version=tlsVersion) if args.insecure: mqttc.tls_insecure_set(True) if args.username or args.password: mqttc.username_pw_set(args.username, args.password) mqttc.on_message = on_message mqttc.on_connect = on_connect mqttc.on_publish = on_publish mqttc.on_subscribe = on_subscribe if args.debug: mqttc.on_log = on_log print("Connecting to "+args.host+" port: "+str(port)) mqttc.connect(args.host, port, args.keepalive) mqttc.subscribe(args.topic, args.qos) mqttc.loop_forever()
提问
请问有人知道该异常的原因是什么吗?
内容的提问来源于stack exchange,提问作者DrBwts
相关产品推荐
相关产品推荐

