7.2.3 MQTT 协议
2026/7/18大约 6 分钟常用组件组件网络MQTTIoT
7.2.3 MQTT 协议
📚 本节导读
学习时长:约 35 分钟
难度级别:⭐⭐⭐⭐☆
前置知识:MQTT 协议基础、发布/订阅模型、第 7.2.1 节 Socket 基础
🎯 学习目标
- 理解 MQTT 协议的发布/订阅模型和 QoS 机制
- 掌握 Paho MQTT 客户端的使用(连接、发布、订阅、心跳)
- 掌握 Mini MQTT 轻量级客户端的使用
- 了解遗嘱消息(Last Will)的作用
- 掌握 MQTT 主题设计最佳实践
一、概述
MQTT(Message Queuing Telemetry Transport)是一种基于发布/订阅模式的轻量级物联网通信协议,具有报文小、带宽占用低、支持不稳定网络等特点,是物联网领域最广泛使用的协议之一。
OneOS 提供两种 MQTT 客户端实现:
| 实现 | 路径 | 特点 | 适用场景 |
|---|---|---|---|
| Paho MQTT | components/net/protocols/mqtt/pahomqtt-v1.1.0/ | 功能完整,支持 QoS 0/1/2,遗嘱消息,持久会话 | 功能需求丰富的物联网设备 |
| Mini MQTT | components/net/protocols/mqtt/minimqtt/ | 代码精简,资源占用极小,API 简洁 | 资源极度受限的设备 |
二、MQTT 协议核心概念
2.1 发布/订阅模型
┌──────────┐ publish "sensor/temp" ┌──────────────┐
│ Publisher │ ─────────────────────────────► │ MQTT Broker │
│ (传感器) │ │ (消息中间件) │
└──────────┘ └──────┬───────┘
│
subscribe "sensor/temp"
│
┌────────▼──────┐
│ Subscriber │
│ (监控中心) │
└───────────────┘MQTT 协议的核心角色:
- Publisher(发布者):向指定主题发布消息
- Subscriber(订阅者):订阅感兴趣的主题,接收消息
- Broker(代理/服务器):负责消息路由和分发
- Topic(主题):消息的分类标识,支持层级通配符(
+匹配单层,#匹配多层)
2.2 QoS 服务质量
| QoS | 名称 | 描述 | 消息语义 |
|---|---|---|---|
| 0 | 最多一次(At most once) | 消息只发送一次,不确认,不重试 | 可能丢失 |
| 1 | 至少一次(At least once) | 消息保证送达,需要确认,可能重复 | 可能重复 |
| 2 | 恰好一次(Exactly once) | 消息保证恰好送达一次,四步握手 | 不丢不重 |
QoS 选择建议:
- QoS 0:适用于传感器高频上报、允许偶尔丢失的场景
- QoS 1:适用于大多数物联网场景,确保消息送达
- QoS 2:适用于关键控制指令,如远程开关机、固件升级触发等
三、Paho MQTT 详细说明
3.1 核心 API 参考
/* 初始化 MQTT 客户端 */
int MQTTClient_init(MQTTClient *client);
/* 连接到 MQTT Broker */
int MQTTClient_connect(MQTTClient *client, MQTTPacket_connectData *options);
/* 订阅主题 */
int MQTTClient_subscribe(MQTTClient *client, const char *topic,
enum QoS qos, messageHandler handler);
/* 发布消息 */
int MQTTClient_publish(MQTTClient *client, const char *topic,
int payloadlen, void *payload,
enum QoS qos, int retained);
/* 取消订阅 */
int MQTTClient_unsubscribe(MQTTClient *client, const char *topic);
/* 断开连接 */
int MQTTClient_disconnect(MQTTClient *client);
/* 循环处理(必须周期性调用) */
int MQTTClient_cycle(MQTTClient *client, int timeout_ms);
/* 检查客户端是否已连接 */
int MQTTClient_is_connected(MQTTClient *client);3.2 连接选项结构体
typedef struct {
const char *clientID; /* 客户端唯一标识 */
const char *username; /* 用户名 */
const char *password; /* 密码 */
const char *willTopic; /* 遗嘱主题 */
const char *willMessage; /* 遗嘱消息内容 */
enum QoS willQoS; /* 遗嘱消息 QoS */
int willRetained; /* 遗嘱消息是否保留 */
int keepAliveInterval; /* 心跳间隔(秒),建议 30-120 秒 */
int cleanSession; /* 是否清理会话 */
int MQTTVersion; /* MQTT 协议版本:3 = v3.1,4 = v3.1.1 */
} MQTTPacket_connectData;遗嘱消息(Last Will and Testament):当客户端非正常断开连接时,Broker 自动将遗嘱消息发布到指定主题,用于检测设备离线状态。
3.3 Paho MQTT 完整示例
#include <MQTTClient.h>
#include <oneos/os_kernel.h>
static MQTTClient g_mqtt_client;
static char g_recv_buf[1024];
static char g_send_buf[1024];
static void message_arrived(MessageData *data)
{
os_kprintf("[MQTT] topic: %.*s\r\n",
data->topicName->lenstring.len,
data->topicName->lenstring.data);
os_kprintf("[MQTT] payload: %.*s\r\n",
data->message->payloadlen,
(char *)data->message->payload);
os_kprintf("[MQTT] qos=%d, retained=%d\r\n",
data->message->qos, data->message->retained);
}
void mqtt_example_task(void *arg)
{
int ret;
MQTTPacket_connectData conn_data = {
.clientID = "oneos_device_001",
.keepAliveInterval = 60,
.cleanSession = 1,
.MQTTVersion = 4,
.willTopic = "device/001/status",
.willMessage = "offline",
.willQoS = QOS1,
.willRetained = 1,
};
ret = MQTTClient_init(&g_mqtt_client);
if (ret != 0) {
os_kprintf("[MQTT] init failed: %d\r\n", ret);
return;
}
os_kprintf("[MQTT] connecting to broker...\r\n");
ret = MQTTClient_connect(&g_mqtt_client, &conn_data);
if (ret != 0) {
os_kprintf("[MQTT] connect failed: %d\r\n", ret);
return;
}
os_kprintf("[MQTT] connected successfully!\r\n");
MQTTClient_publish(&g_mqtt_client, "device/001/status",
strlen("online"), "online", QOS1, 1);
ret = MQTTClient_subscribe(&g_mqtt_client, "device/001/cmd",
QOS1, message_arrived);
if (ret != 0) {
os_kprintf("[MQTT] subscribe failed: %d\r\n", ret);
} else {
os_kprintf("[MQTT] subscribed to 'device/001/cmd'\r\n");
}
MQTTClient_subscribe(&g_mqtt_client, "device/+/broadcast",
QOS1, message_arrived);
int counter = 0;
while (1) {
char payload[64];
snprintf(payload, sizeof(payload),
"{\"temp\":25.5,\"humi\":60.0,\"seq\":%d}", counter++);
MQTTClient_publish(&g_mqtt_client, "device/001/data",
strlen(payload), payload, QOS1, 0);
ret = MQTTClient_cycle(&g_mqtt_client, 1000);
if (ret != 0) {
os_kprintf("[MQTT] cycle error: %d, reconnecting...\r\n", ret);
MQTTClient_disconnect(&g_mqtt_client);
os_task_tsleep(2000);
ret = MQTTClient_connect(&g_mqtt_client, &conn_data);
if (ret == 0) {
MQTTClient_subscribe(&g_mqtt_client, "device/001/cmd",
QOS1, message_arrived);
MQTTClient_subscribe(&g_mqtt_client, "device/+/broadcast",
QOS1, message_arrived);
MQTTClient_publish(&g_mqtt_client, "device/001/status",
strlen("online"), "online", QOS1, 1);
}
}
os_task_tsleep(5000);
}
}四、Mini MQTT
4.1 核心 API 参考
#include <m_mqtt_client.h>
/* 连接 MQTT Broker */
int mqtt_connect(m_mqtt_client_t *client, const char *host, int port,
const char *client_id, const char *user, const char *pass);
/* 发布消息 */
int mqtt_publish(m_mqtt_client_t *client, const char *topic,
const char *payload, int len, int qos, int retain);
/* 订阅主题 */
int mqtt_subscribe(m_mqtt_client_t *client, const char *topic, int qos);
/* 循环处理 */
int mqtt_loop(m_mqtt_client_t *client, int timeout_ms);
/* 断开连接 */
int mqtt_disconnect(m_mqtt_client_t *client);4.2 Mini MQTT 使用示例
#include <m_mqtt_client.h>
#include <oneos/os_kernel.h>
static m_mqtt_client_t g_client;
void minimqtt_example(void *arg)
{
int ret;
ret = mqtt_connect(&g_client, "broker.emqx.io", 1883,
"oneos_mini_001", NULL, NULL);
if (ret != 0) {
os_kprintf("[MiniMQTT] connect failed: %d\r\n", ret);
return;
}
os_kprintf("[MiniMQTT] connected!\r\n");
mqtt_subscribe(&g_client, "test/topic", 1);
mqtt_publish(&g_client, "test/topic", "Hello MiniMQTT", 14, 1, 0);
while (1) {
ret = mqtt_loop(&g_client, 1000);
if (ret > 0) {
os_kprintf("[MiniMQTT] received message\r\n");
} else if (ret < 0) {
os_kprintf("[MiniMQTT] connection lost\r\n");
break;
}
os_task_tsleep(500);
}
mqtt_disconnect(&g_client);
}五、MQTT 主题设计最佳实践
主题命名规范:
{产品}/{设备ID}/{功能}
示例:
sensor/temp001/data -- 温度传感器数据上报
sensor/temp001/cmd -- 温度传感器命令下发
gateway/001/status -- 网关状态
device/+/event -- 所有设备的事件(单层通配符)
device/# -- 所有设备的所有消息(多层通配符)设计原则:
- 主题层级不宜过深(建议不超过 5 层)
- 不要以
/开头或结尾 - 使用可读的层级名称,避免缩写混淆
- 区分数据上报(data)和命令下发(cmd)主题
- 为每个设备预留状态主题(status),配合遗嘱消息实现在线/离线检测
📝 本节小结
本节介绍了 OneOS 中 MQTT 协议的使用。OneOS 提供 Paho MQTT(功能完整)和 Mini MQTT(轻量精简)两种客户端实现,支持 QoS 0/1/2、遗嘱消息、持久会话等核心特性。
核心要点回顾:
- MQTT 基于发布/订阅模型,通过 Broker 进行消息路由
- Paho MQTT 功能完整,适合大多数物联网设备
- Mini MQTT 代码精简,适合资源极度受限的设备
MQTTClient_cycle/mqtt_loop必须周期性调用以维持心跳和接收消息- 遗嘱消息可检测设备离线状态,是物联网场景的重要机制