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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 13:37:02