以下是关于“QoS 1 消息传输:Mosquitto C/C++/Python 客户端实测”的详细回答。我将逐步解释QoS 1的概念,并提供在Mosquitto消息代理环境下,使用C、C++和Python客户端的实测代码示例。测试基于真实MQTT协议行为,确保回答可靠。实测环境建议:

  • 安装Mosquitto broker(如Ubuntu下 sudo apt-get install mosquitto)。
  • 安装客户端库:Python使用 paho-mqtt,C/C++使用Paho MQTT库(下载地址:Eclipse Paho)。
  • 测试方法:启动broker,运行发送方和接收方客户端,观察日志(如消息重传、PUBACK确认)来验证QoS 1行为(消息至少一次传递,可能重复)。

1. QoS 1 背景知识

在MQTT协议中,QoS(服务质量)级别1确保消息至少被传递一次。发送方发布消息后,等待接收方的PUBACK确认;如果未收到确认,发送方会重传消息。这适用于可靠性要求较高的场景,如传感器数据传输。Mosquitto是一个开源MQTT broker,支持多种客户端库。

2. 实测准备

  • 环境设置:启动Mosquitto broker(默认端口1883)。命令:mosquitto -v(启用详细日志)。
  • 测试场景:发送方发布消息到主题(如 test/topic),接收方订阅同一主题。模拟网络不稳定(如断开连接)来观察QoS 1重传行为。
  • 验证指标
    • 发送方:日志显示消息发布和重传。
    • 接收方:日志显示消息接收和PUBACK发送。
    • 使用工具:如 mosquitto_submosquitto_pub 命令行工具辅助测试。

3. Python 客户端实测

使用 paho-mqtt 库(安装:pip install paho-mqtt)。以下代码分为发送方和接收方,设置QoS=1。

发送方代码(发布消息)
import paho.mqtt.client as mqtt
import time

def on_publish(client, userdata, mid):
    print(f"消息已发布, MID: {mid}")  # MID是消息ID,用于跟踪QoS 1确认

client = mqtt.Client()
client.on_publish = on_publish
client.connect("localhost", 1883, 60)
client.loop_start()

# 发布QoS 1消息
topic = "test/topic"
message = "Hello QoS 1"
client.publish(topic, message, qos=1)
print("消息发布中...")

time.sleep(2)  # 等待确认
client.disconnect()

接收方代码(订阅消息)
import paho.mqtt.client as mqtt

def on_connect(client, userdata, flags, rc):
    print("连接成功,订阅主题")
    client.subscribe("test/topic", qos=1)  # QoS 1订阅

def on_message(client, userdata, msg):
    print(f"收到消息: {msg.payload.decode()}, QoS: {msg.qos}")
    # QoS 1下,接收方自动发送PUBACK

client = mqtt.Client()
client.on_connect = on_connect
client.on_message = on_message
client.connect("localhost", 1883, 60)
client.loop_forever()  # 持续监听

测试步骤
  1. 启动Mosquitto broker。
  2. 运行接收方代码。
  3. 运行发送方代码。
  4. 观察行为
    • 正常网络:接收方打印消息一次,发送方显示 消息已发布
    • 模拟失败:断开接收方网络后重连,发送方日志显示重传(MID不变),接收方可能收到重复消息。

4. C 客户端实测

使用Paho MQTT C库(安装:下载源码编译)。代码基于异步模式,设置QoS=1。

发送方代码
#include <stdio.h>
#include <stdlib.h>
#include <MQTTClient.h>

#define ADDRESS "tcp://localhost:1883"
#define CLIENTID "C_Sender"
#define TOPIC "test/topic"
#define QOS 1

void delivered(void *context, MQTTClient_deliveryToken dt) {
    printf("消息已交付, Token: %d\n", dt);
}

int main() {
    MQTTClient client;
    MQTTClient_create(&client, ADDRESS, CLIENTID, MQTTCLIENT_PERSISTENCE_NONE, NULL);
    MQTTClient_connectOptions conn_opts = MQTTClient_connectOptions_initializer;
    MQTTClient_setCallbacks(client, NULL, NULL, delivered, NULL); // 设置交付回调

    if (MQTTClient_connect(client, &conn_opts) != MQTTCLIENT_SUCCESS) {
        fprintf(stderr, "连接失败\n");
        exit(EXIT_FAILURE);
    }

    char *payload = "Hello QoS 1 from C";
    MQTTClient_message pubmsg = MQTTClient_message_initializer;
    pubmsg.payload = payload;
    pubmsg.payloadlen = strlen(payload);
    pubmsg.qos = QOS;
    pubmsg.retained = 0;

    MQTTClient_deliveryToken token;
    MQTTClient_publishMessage(client, TOPIC, &pubmsg, &token);
    printf("发布消息中...\n");
    MQTTClient_waitForCompletion(client, token, 1000L); // 等待确认

    MQTTClient_disconnect(client, 1000);
    MQTTClient_destroy(&client);
    return 0;
}

接收方代码
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <MQTTClient.h>

#define ADDRESS "tcp://localhost:1883"
#define CLIENTID "C_Receiver"
#define TOPIC "test/topic"
#define QOS 1

int msgarrvd(void *context, char *topicName, int topicLen, MQTTClient_message *message) {
    printf("收到消息: %.*s, QoS: %d\n", message->payloadlen, (char*)message->payload, message->qos);
    MQTTClient_freeMessage(&message); // 释放内存
    MQTTClient_free(topicName);
    return 1; // 返回1表示处理成功,自动发送PUBACK
}

int main() {
    MQTTClient client;
    MQTTClient_create(&client, ADDRESS, CLIENTID, MQTTCLIENT_PERSISTENCE_NONE, NULL);
    MQTTClient_connectOptions conn_opts = MQTTClient_connectOptions_initializer;
    MQTTClient_setCallbacks(client, NULL, NULL, msgarrvd, NULL); // 设置消息回调

    if (MQTTClient_connect(client, &conn_opts) != MQTTCLIENT_SUCCESS) {
        fprintf(stderr, "连接失败\n");
        exit(EXIT_FAILURE);
    }

    MQTTClient_subscribe(client, TOPIC, QOS);
    printf("订阅中...按Enter退出\n");
    getchar(); // 阻塞等待消息

    MQTTClient_unsubscribe(client, TOPIC);
    MQTTClient_disconnect(client, 1000);
    MQTTClient_destroy(&client);
    return 0;
}

测试步骤
  1. 编译代码:使用 gcc sender.c -o sender -lpaho-mqtt3cgcc receiver.c -o receiver -lpaho-mqtt3c
  2. 启动broker。
  3. 运行接收方:./receiver
  4. 运行发送方:./sender
  5. 观察行为:接收方输出消息内容,发送方显示交付令牌。模拟网络中断时,发送方重试(日志显示多次交付回调)。

5. C++ 客户端实测

使用Paho MQTT C++库(基于C库封装)。代码设置QoS=1。

发送方代码
#include <iostream>
#include <cstring>
#include "mqtt/async_client.h"

const std::string SERVER_ADDRESS("tcp://localhost:1883");
const std::string CLIENT_ID("Cpp_Sender");
const std::string TOPIC("test/topic");
const int QOS = 1;

int main() {
    mqtt::async_client client(SERVER_ADDRESS, CLIENT_ID);
    mqtt::connect_options connOpts;
    connOpts.set_clean_session(true);

    try {
        client.connect(connOpts)->wait();
        auto pubmsg = mqtt::make_message(TOPIC, "Hello QoS 1 from C++", QOS, false);
        auto token = client.publish(pubmsg);
        token->wait(); // 等待发布完成
        std::cout << "消息发布成功" << std::endl;
        client.disconnect()->wait();
    } catch (const mqtt::exception& exc) {
        std::cerr << "错误: " << exc.what() << std::endl;
        return 1;
    }
    return 0;
}

接收方代码
#include <iostream>
#include "mqtt/async_client.h"

const std::string SERVER_ADDRESS("tcp://localhost:1883");
const std::string CLIENT_ID("Cpp_Receiver");
const std::string TOPIC("test/topic");
const int QOS = 1;

class callback : public virtual mqtt::callback {
public:
    void message_arrived(mqtt::const_message_ptr msg) override {
        std::cout << "收到消息: " << msg->to_string() << ", QoS: " << msg->get_qos() << std::endl;
    }
};

int main() {
    mqtt::async_client client(SERVER_ADDRESS, CLIENT_ID);
    callback cb;
    client.set_callback(cb);

    mqtt::connect_options connOpts;
    connOpts.set_clean_session(true);

    try {
        client.connect(connOpts)->wait();
        client.subscribe(TOPIC, QOS)->wait();
        std::cout << "订阅中...按Enter退出" << std::endl;
        std::cin.get(); // 阻塞等待
        client.unsubscribe(TOPIC)->wait();
        client.disconnect()->wait();
    } catch (const mqtt::exception& exc) {
        std::cerr << "错误: " << exc.what() << std::endl;
        return 1;
    }
    return 0;
}

测试步骤
  1. 编译:需要链接Paho C++库(参考官方文档)。
  2. 启动broker。
  3. 运行接收方。
  4. 运行发送方。
  5. 观察行为:类似C客户端,但C++更面向对象。QoS 1下,重复消息可通过接收方日志识别。

6. 实测总结

  • QoS 1 行为验证:在所有语言中,当网络稳定时,消息传递一次;当模拟丢包(如kill接收方进程后重启),发送方重传消息,接收方可能收到重复(需应用层去重)。
  • 性能注意:QoS 1增加延迟和带宽使用(因确认机制),在高频场景测试时监控资源。
  • 建议:使用Wireshark抓包或Mosquitto日志(-v 参数)详细跟踪PUBLISH和PUBACK包。实测代码基于Paho库标准用法,确保可靠。如需扩展,参考 Eclipse Paho文档

如果您有特定场景(如大规模测试或错误处理),请提供更多细节,我可以进一步优化代码!

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐