1. 背景
1.1 MQTT 是什么
MQTT(Message Queuing Telemetry Transport,消息队列遥测传输)是一种基于发布/订阅(Pub/Sub)模式的轻量级消息传输协议。它由 IBM 的 Andy Stanford-Clark 与 Arcom 的 Arlen Nipper 于 1999 年发明,最初用于石油天然气管道遥测;2014 年成为 OASIS 开放标准,2016 年纳入 ISO/IEC 20922 国际标准。
三个核心设计动机:
| 设计动机 | 具体体现 |
|---|---|
| 极轻量 | 固定报文头最小仅 2 字节(CONNECT/PUBLISH 等报文),非常适合低带宽、高延迟、不稳定的网络环境 |
| 发布/订阅解耦 | 发布者与订阅者互不感知,通过 Broker(消息代理)中转,天然支持一对多广播与多对一汇聚 |
| QoS 可靠性分级 | 提供 0 / 1 / 2 三级服务质量,在"带宽占用"与"送达可靠性"之间可权衡 |
1.2 MQTT 3.1.1 与 MQTT 5.0
- MQTT 3.1.1(2014,最普及):固定头 + 可变头 + Payload,QoS 0/1/2,遗嘱消息(Will Message)、保留消息(Retained Message)、持久会话(Clean Session)等核心能力已齐全。
- MQTT 5.0(2019):新增会话过期(Session Expiry)、用户属性(User Properties)、主题别名(Topic Alias)、共享订阅(Shared Subscription)、请求/响应模式、原因码(Reason Code)等。Paho C 库同时支持两个版本。
1.3 Eclipse Paho 项目
Paho 是 Eclipse 基金会旗下的开源 MQTT 客户端实现家族,覆盖 C/C++、Java、Python、Go、JavaScript 等语言。paho.mqtt.c 是纯 C 语言实现:
- 支持 MQTT 3.1 / 3.1.1 / 5.0,支持 QoS 0/1/2
- 提供两套 API:同步 API(MQTTClient_*)与异步 API(MQTTAsync_*)
- 支持 TLS/SSL(依赖 OpenSSL)、WebSocket 传输
- 采用 Eclipse Public License 2.0 / Eclipse Distribution License 1.0 双许可
- 被大量边缘网关、嵌入式设备、工业数采方案采用
1.4 为什么工业数采场景适合 MQTT
工业数采的典型链路:CNC/PLC 设备 → 现场总线采集(Modbus / OPC UA)→ 边缘网关汇聚 → 消息总线上云 → 平台侧存储分析。MQTT 在这一链路中的位置正是"边缘网关 → 云平台"这一段:
- 采集点数量大(一条产线几十上百台设备),MQTT 的发布/订阅天然支持多网关汇聚到单一 Topic 树;
- 网络环境可能不稳定(车间 Wi-Fi、4G 网关),MQTT 的 QoS 1/2、遗嘱消息、持久会话能保证数据尽量不丢;
- Broker 侧生态成熟(EMQX / Mosquitto / HiveMQ),云平台可直接订阅消费。
2. 核心 API 说明
2.1 同步 API(MQTTClient_*)
同步 API 在独立后台线程中维持网络连接,用户调用阻塞函数进行收发,心智模型简单:
| API | 作用 |
|---|---|
| MQTTClient_create(&client, serverURI, clientId, persistence_type, persistence_context) | 创建客户端句柄;persistence_type 可为 MQTTCLIENT_PERSISTENCE_NONE(内存)或 MQTTCLIENT_PERSISTENCE_DEFAULT(磁盘,重启后可恢复会话) |
| MQTTClient_connect(client, &connOpts) | 建立连接;connOpts 中可配置 keepAlive、cleanSession、遗嘱、用户名密码、超时等 |
| MQTTClient_subscribe(client, topic, qos) | 订阅主题,返回订阅回执 QoS |
| MQTTClient_publishMessage(client, topic, &msg, &token) | 发布消息(内部会拷贝payload,可安全复用缓冲区) |
| MQTTClient_receive(client, &topicName, &topicLen, &msg, timeout_ms) | 阻塞接收消息(超时返回 MQTTCLIENT_SUCCESS 但 topicName 为 NULL 表示无消息);须调用 MQTTClient_freeMessage(&msg) 与 MQTTClient_free(topicName) 释放 |
| MQTTClient_disconnect(client, timeout_ms) | 优雅断开(发送 DISCONNECT,触发遗嘱前清理) |
| MQTTClient_destroy(&client) | 销毁句柄,释放资源 |
回调机制(同步 API 同样支持,由后台线程触发):
| 回调 | 触发时机 |
|---|---|
| connectionLost | 网络连接意外断开(非主动 disconnect) |
| messageArrived | 收到订阅主题的消息(返回 1 表示消息已处理,0 表示交给 receive 队列) |
| deliveryComplete | QoS > 0 的消息送达确认完成 |
关键点:同步 API 的 messageArrived 回调返回 1 时消息被消费;返回 0 时消息进入内部队列,可再通过 MQTTClient_receive 取出。
2.2 异步 API(MQTTAsync_*)
异步 API 完全由回调驱动,不提供阻塞 receive,适合高吞吐、事件驱动架构:
| API | 作用 |
|---|---|
| MQTTAsync_create(&client, serverURI, clientId, persistence_type, persistence_context) | 创建异步客户端 |
| MQTTAsync_setCallbacks(client, context, connectionLost, messageArrived, deliveryComplete) | 注册连接丢失 / 消息到达 / 送达完成回调 |
| MQTTAsync_connect(client, &connOpts) | 异步发起连接,结果通过 onSuccess / onFailure 回调返回 |
| MQTTAsync_subscribe(client, topic, qos, response, onSuccess, onFailure, context) | 异步订阅 |
| MQTTAsync_publishMessage(client, topic, &msg, response, onSuccess, onFailure, context) | 异步发布;payload 在 deliveryComplete 回调前必须保持存活(异步 API 不拷贝 payload) |
| MQTTAsync_disconnect(client, &discOpts) | 异步断开 |
2.3 核心配置结构体
/* 连接选项 */ MQTTClient_connectOptions connOpts = MQTTClient_connectOptions_initializer; connOpts.keepAliveInterval = 20; /* 心跳间隔(秒),0 表示关闭 keepalive */ connOpts.cleansession = 1; /* 1=每次连接清空会话;0=持久会话恢复 */ connOpts.username = "user"; connOpts.password = "pass"; connOpts.connectTimeout = 10; /* 连接超时(秒) */ connOpts.will = &willOpts; /* 遗嘱消息 */ connOpts.ssl = &sslOpts; /* TLS 配置 */ connOpts.retryInterval = 5; /* 重传间隔 */ connOpts.sessionExpiryInterval = 3600; /* MQTT 5.0 会话过期秒数 */ /* 遗嘱消息 */ MQTTClient_willOptions willOpts = MQTTClient_willOptions_initializer; willOpts.topicName = "device/offline"; willOpts.message = "unexpected down"; willOpts.qos = 1; willOpts.retained = 1; /* 消息 */ MQTTClient_message pubmsg = MQTTClient_message_initializer; pubmsg.payload = (void*)payload; pubmsg.payloadlen = (int)strlen(payload); pubmsg.qos = 1; pubmsg.retained = 0;cleanSession(MQTT 5.0 中为 Clean Start + Session Expiry)语义极其重要:
- cleansession = 1:连接成功后立即丢弃旧会话;断线后 Broker 清除该客户端的订阅与离线消息。
- cleansession = 0:持久会话。断线重连(同一 clientId)后,Broker 恢复订阅、补发断线期间的 QoS 1/2 离线消息;前提是会话未过期(MQTT 5.0 受 sessionExpiryInterval 控制)。
2.4 QoS 语义速查
| QoS | 名称 | 保证 | 可能副作用 | 适用 |
|---|---|---|---|---|
| 0 | At most once | 最多一次,尽力而为 | 可能丢失 | 传感器高频采样、允许丢点 |
| 1 | At least once | 至少一次 | 可能重复 | 设备状态、命令下发(配合幂等) |
| 2 | Exactly once | 恰好一次 | 开销最大 | 计费、告警、需精确记账 |
3. 详细使用说明
3.1 安装与 CMake 集成
vcpkg:
vcpkg install paho-mqttCMake FetchContent:
include(FetchContent) FetchContent_Declare(paho_mqtt_c GIT_REPOSITORY https://github.com/eclipse-paho/paho.mqtt.c.git GIT_TAG v1.3.13) FetchContent_MakeAvailable(paho_mqtt_c)find_package(vcpkg 安装后):
find_package(PahoMqttC REQUIRED) target_link_libraries(my_app PRIVATE PahoMqttC::paho-mqtt3c) # 同步 API # 或 PahoMqttC::paho-mqtt3cs(同步+SSL)、paho-mqtt3a(异步)、paho-mqtt3as(异步+SSL)3.2 同步 API 最小完整示例
#include "MQTTClient.h" #include <stdio.h> #include <string.h> #define ADDRESS "tcp://localhost:1883" #define CLIENTID "gateway-01" #define TOPIC "factory/line1/cnc01/status" #define PAYLOAD "{\"state\":\"running\",\"spindle_rpm\":8000}" #define QOS 1 volatile int msg_arrived = 0; int on_message(void *ctx, char *topicName, int topicLen, MQTTClient_message *message) { printf("recv [%s]: %.*s\n", topicName, message->payloadlen, (char*)message->payload); MQTTClient_freeMessage(&message); MQTTClient_free(topicName); msg_arrived = 1; return 1; /* 已消费,不进入 receive 队列 */ } int main(void) { MQTTClient client; MQTTClient_connectOptions connOpts = MQTTClient_connectOptions_initializer; MQTTClient_message pubmsg = MQTTClient_message_initializer; MQTTClient_deliveryToken token; int rc; MQTTClient_create(&client, ADDRESS, CLIENTID, MQTTCLIENT_PERSISTENCE_NONE, NULL); connOpts.keepAliveInterval = 20; connOpts.cleansession = 1; connOpts.connectTimeout = 10; MQTTClient_setCallbacks(client, NULL, NULL, on_message, NULL); if ((rc = MQTTClient_connect(client, &connOpts)) != MQTTCLIENT_SUCCESS) { printf("connect failed: %d\n", rc); return 1; } printf("connected to %s\n", ADDRESS); MQTTClient_subscribe(client, TOPIC, QOS); pubmsg.payload = (void*)PAYLOAD; pubmsg.payloadlen = (int)strlen(PAYLOAD); pubmsg.qos = QOS; pubmsg.retained = 0; MQTTClient_publishMessage(client, TOPIC, &pubmsg, &token); printf("published, token=%d\n", token); /* 等待回调触发(演示用;实际生产用事件循环/信号量) */ while (!msg_arrived) { /* 简单忙等或 sleep */ } MQTTClient_disconnect(client, 1000); MQTTClient_destroy(&client); return 0; }编译(CMake 或直接):
gcc demo.c -o demo -lpaho-mqtt3c3.3 异步 API 最小示例(回调驱动)
#include "MQTTAsync.h" #include <stdio.h> #include <string.h> #define ADDRESS "tcp://localhost:1883" #define CLIENTID "async-gateway" #define TOPIC "factory/line1/cnc01/status" #define QOS 1 static volatile int finished = 0; static void on_connect_success(void *ctx, MQTTAsync_successData *resp) { printf("connect ok\n"); } static void on_connect_failure(void *ctx, MQTTAsync_failureData *resp) { printf("connect fail rc=%d\n", resp ? resp->code : -1); finished = 1; } static int on_message(void *ctx, char *topic, int topicLen, MQTTAsync_message *msg) { printf("recv [%s]: %.*s\n", topic, msg->payloadlen, (char*)msg->payload); MQTTAsync_freeMessage(&msg); MQTTAsync_free(topic); finished = 1; return 1; } static void on_publish_success(void *ctx, MQTTAsync_successData *resp) { printf("publish delivered\n"); } int main(void) { MQTTAsync client; MQTTAsync_connectOptions connOpts = MQTTAsync_connectOptions_initializer; MQTTAsync_message pubmsg = MQTTAsync_message_initializer; MQTTAsync_create(&client, ADDRESS, CLIENTID, MQTTCLIENT_PERSISTENCE_NONE, NULL); MQTTAsync_setCallbacks(client, NULL, NULL, on_message, NULL); connOpts.keepAliveInterval = 20; connOpts.cleansession = 1; connOpts.onSuccess = on_connect_success; connOpts.onFailure = on_connect_failure; MQTTAsync_connect(client, &connOpts); /* 等待 on_connect_success 后再发布(生产用状态机/信号量同步) */ MQTTAsync_subscribe(client, TOPIC, QOS, NULL, NULL, NULL, NULL); pubmsg.payload = "async hello"; pubmsg.payloadlen = 11; pubmsg.qos = QOS; MQTTAsync_publishMessage(client, TOPIC, &pubmsg, NULL, on_publish_success, NULL, NULL); /* 注意:异步 API 不拷贝 payload,deliveryComplete/onSuccess 前必须保持存活 */ while (!finished) { /* 等消息 */ } MQTTAsync_disconnect(client, NULL); MQTTAsync_destroy(&client); return 0; }3.4 遗嘱消息与保留消息
遗嘱(Last Will and Testament, LWT):连接建立时登记遗嘱主题与消息。当 Broker 检测到客户端非正常断开(网络中断、心跳超时、崩溃)时,代为发布遗嘱消息;客户端主动 disconnect 时不触发。
典型用法:设备上线发布 device/online,连接时登记遗嘱 device/offline,平台通过在线状态判断设备存活。
MQTTClient_willOptions willOpts = MQTTClient_willOptions_initializer; willOpts.topicName = "factory/line1/cnc01/offline"; willOpts.message = "{\"state\":\"down\"}"; willOpts.qos = 1; willOpts.retained = 1; connOpts.will = &willOpts;保留消息(Retained Message):发布时置 retained = 1,Broker 会保存该主题的最后一条消息,新订阅者订阅该主题时立即收到。适合"设备当前状态""配置信息"这类"只关心最新值"的数据。
3.5 断线重连
Paho 同步 API 的 MQTTClient_connect 失败后需自行实现重连逻辑;connectionLost 回调中通常做重连。常用模式:
int on_connection_lost(void *ctx, char *cause) { /* 记录日志、进入重连状态机 */ return 1; } /* 主循环:连接失败/断开后按退避策略重连 */ while (running) { rc = MQTTClient_connect(client, &connOpts); if (rc == MQTTCLIENT_SUCCESS) break; sleep(backoff); /* 建议指数退避 + 抖动 */ }配合 cleansession = 0 时,重连后 Broker 会补发断线期间的 QoS 1/2 离线消息。
3.6 TLS 连接
MQTTClient_SSLOptions sslOpts = MQTTClient_SSLOptions_initializer; sslOpts.trustStore = "ca.crt"; /* CA 证书;NULL 表示使用系统默认 */ sslOpts.keyStore = "client.pem"; /* 客户端证书(双向认证时) */ sslOpts.privateKey = "client.key"; sslOpts.enableServerCertAuth = 1; /* 校验服务器证书 */ connOpts.ssl = &sslOpts; /* 地址改为 ssl://host:8883 */4. 常错点 / 坑(20 条)
| 坑 | 说明与规避 | |
|---|---|---|
| 1 | 同步 API receive 忘释放 | MQTTClient_receive 返回的消息必须 MQTTClient_freeMessage + MQTTClient_free(topicName),否则内存泄漏 |
| 2 | 异步 API payload 被提前释放 | 异步 publish不拷贝payload,需保证缓冲区存活到 deliveryComplete/onSuccess;同步 API 才拷贝 |
| 3 | cleansession 语义误解 | =0 表示持久会话:重连恢复订阅并补发离线消息;=1 表示每次全新开始。乱用导致重复消息或"消息丢了"的错觉 |
| 4 | 遗嘱触发条件搞错 | 仅非正常断开触发;主动 MQTTClient_disconnect 不会发遗嘱。测试时常用 kill -9 模拟 |
| 5 | QoS 1 重复投递未做幂等 | QoS 1 是 at-least-once,可能重复;消费端需按消息 ID/时间戳/设备状态幂等去重 |
| 6 | 忘记处理 messageArrived 返回语义 | 返回 1 消费、返回 0 进 receive 队列;混用两套消费方式会丢消息或重复处理 |
| 7 | 连接后不检查 MQTTClient_connect 返回值 | 连接失败(网络不通、认证失败、clientId 冲突)仍继续 publish,表现为"静默丢数据" |
| 8 | keepAliveInterval 设置不当 | 过短(<5s)造成频繁心跳浪费带宽;过长(>120s)导致断线检测迟钝;防火墙/NAT 环境下建议 20~60s |
| 9 | 心跳被 NAT/防火墙掐断未察觉 | 长连接场景建议 keepAlive 小于 NAT 超时;或使用 MQTT 5.0 的 Keep Alive + 服务器侧断连检测 |
| 10 | 主题通配符误用 | 订阅可含 +(单层)、#(多层,必须放最后);发布不允许通配符。a/# 会匹配 a、a/b、a/b/c |
| 11 | 订阅回执 QoS 与请求不符 | Broker 可降级订阅 QoS(请求 2 可能回 1);需读取 subscribe 返回值,不要假设与请求一致 |
| 12 | 多线程共享同一客户端 | 同一 MQTTClient 句柄并发 publish 需自行加锁,或使用异步 API + 串行发布;Paho 内部队列非完全线程安全 |
| 13 | clientId 冲突 | 两个客户端用相同 clientId 连接同一 Broker,后连者会把先连者踢下线(协议规定),表现为"连接被随机断开" |
| 14 | 忘配置 connectTimeout | 默认值可能过长,断网时 connect 阻塞很久;显式设置 5~10s |
| 15 | 大 payload 缓冲假设 | 默认消息最大约 256MB(受 MAX_MSG_SIZE 限制),超大 payload 需确认 Broker 端 max_packet_size 配置 |
| 16 | 中文/UTF-8 主题 | MQTT 主题和 payload 都是 UTF-8 字节流;Windows 下用窄字符/GBK 发送中文 topic 会导致 Broker 拒绝(topic 必须合法 UTF-8) |
| 17 | TLS 自签名证书信任 | enableServerCertAuth=1 且未把自签名 CA 加入 trustStore 会握手失败;测试环境常用 enableServerCertAuth=0,生产禁止 |
| 18 | 忽略返回码含义 | MQTTCLIENT_SUCCESS=0、MQTTCLIENT_DISCONNECTED=-3、MQTTCLIENT_TOPICNAME_TRUNCATED=-7 等;统一用 MQTTClient_strerror(rc) 打日志 |
| 19 | 断线重连无退避 | 连接失败立即重试会造成"重连风暴"(Broker 端踢号、日志刷屏);用指数退避 + 随机抖动 |
| 20 | MQTT 5.0 属性 API 误用 | 5.0 的 MQTTProperties 需初始化 MQTTProperties_initializer、MQTTProperties_add 添加、用完 MQTTProperties_free;忘了会泄漏或读到脏数据 |
5. 总结
5.1 适用场景表
| 场景 | 推荐 |
|---|---|
| 边缘网关 → 云平台数据上送(工业数采主链路) | MQTT QoS 1 + 遗嘱 + 持久会话 |
| 设备在线状态监控 | 保留消息(最后状态)+ 遗嘱(离线通知) |
| 高频传感器采样(允许丢点) | QoS 0,降低带宽与 Broker 压力 |
| 命令下发(需确认) | QoS 1/2 + 业务层幂等 |
| 海量日志流、分区顺序消费 | 选 Kafka(librdkafka / kafka-go),不是 MQTT |
5.2 选型决策树
机器间通信需求? ├─ 同机进程间共享数据 → Boost.Interprocess ├─ 跨机器轻量遥测/设备上云 → MQTT(本文) │ └─ 需要语义建模/信息模型 → OPC UA(open62541) ├─ 海量日志/事件流、分区有序 → Kafka(librdkafka / kafka-go) ├─ 去中心化、无 Broker → ZeroMQ └─ RPC 调用、跨语言服务 → gRPC
5.3 工业数采实践建议
结合你维护的数采项目(kafka 采集、IPQC 采集、fanuc/mitsubishi 直连):
- 采集分层:CNC/PLC 侧用 libmodbus / open62541 直连采集;边缘网关汇聚后通过 Paho MQTT 上送平台,与 Kafka 链路形成"近场轻量 + 云端海量"双通道。
- Topic 规划:建议 factory/{产线}/{设备类型}/{设备ID}/{指标} 层级,如 factory/line1/cnc/01/spindle_rpm,便于 Broker 权限控制与平台按前缀订阅。
- 可靠性组合拳:QoS 1 + 持久会话(cleansession=0)+ 遗嘱(离线告警)+ 消费端按 device_id + timestamp 幂等去重。
- 安全:生产环境 TLS(Basic256Sha256 级别的双向认证)+ Broker 端 ACL 限制发布前缀。
- 与既有链路衔接:C++ 采集器可用 Paho C(或 C++ 的 paho.mqtt.cpp)发布,Go 网关可用 eclipse/paho.mqtt.golang,与 kafka-go 篇的 Writer/Reader 模式形成跨语言对照。
5.4 FAQ 速查表
| 问题 | 答案 |
|---|---|
| 同步还是异步 API? | 简单采集网关用同步(心智简单);高吞吐/事件驱动用异步 |
| 消息丢了怎么办? | 升级 QoS、持久会话、检查 Broker max_qos/retain_available 配置 |
| 收不到订阅消息? | 查 Topic 拼写/通配符、QoS 是否被降级、是否已订阅成功、clientId 是否被踢 |
| 重复消息如何根治? | QoS 1 天然可能重复,业务层按业务主键幂等;要严格不重选 QoS 2(开销大) |
| 连接总被断开? | 查 clientId 冲突、keepAlive 过短、Broker 连接数上限、心跳是否被 NAT 掐断 |
| 自签 TLS 连不上? | 把自签 CA 放进 trustStore 并 enableServerCertAuth=1,勿直接关校验 |