QoS 1 消息传输:Mosquitto C/C++/Python 客户端实测
·
以下是关于“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_sub和mosquitto_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() # 持续监听
测试步骤
- 启动Mosquitto broker。
- 运行接收方代码。
- 运行发送方代码。
- 观察行为:
- 正常网络:接收方打印消息一次,发送方显示
消息已发布。 - 模拟失败:断开接收方网络后重连,发送方日志显示重传(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;
}
测试步骤
- 编译代码:使用
gcc sender.c -o sender -lpaho-mqtt3c和gcc receiver.c -o receiver -lpaho-mqtt3c。 - 启动broker。
- 运行接收方:
./receiver。 - 运行发送方:
./sender。 - 观察行为:接收方输出消息内容,发送方显示交付令牌。模拟网络中断时,发送方重试(日志显示多次交付回调)。
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;
}
测试步骤
- 编译:需要链接Paho C++库(参考官方文档)。
- 启动broker。
- 运行接收方。
- 运行发送方。
- 观察行为:类似C客户端,但C++更面向对象。QoS 1下,重复消息可通过接收方日志识别。
6. 实测总结
- QoS 1 行为验证:在所有语言中,当网络稳定时,消息传递一次;当模拟丢包(如kill接收方进程后重启),发送方重传消息,接收方可能收到重复(需应用层去重)。
- 性能注意:QoS 1增加延迟和带宽使用(因确认机制),在高频场景测试时监控资源。
- 建议:使用Wireshark抓包或Mosquitto日志(
-v参数)详细跟踪PUBLISH和PUBACK包。实测代码基于Paho库标准用法,确保可靠。如需扩展,参考 Eclipse Paho文档。
如果您有特定场景(如大规模测试或错误处理),请提供更多细节,我可以进一步优化代码!
更多推荐

所有评论(0)