梳理与编写MQTT 客户端程序的通用框架和思路
一、代码总体功能
这是一个运行在嵌入式 Linux 设备上的 MQTT 客户端,主要功能:
-
从配置文件
/etc/wechat.cfg读取连接参数(Broker URI、ClientID、用户名、密码、产品 ID、设备名称)。 -
初始化 RPC 客户端(用于调用板载硬件接口,如读取温湿度传感器、控制 LED)。
-
连接 MQTT 服务器(Broker)。
-
订阅下行主题(接收云端或手机端下发的控制命令,例如开灯/关灯)。
-
循环读取温湿度传感器,构造 JSON 消息,发布到上行主题(上报设备数据)。
-
在回调函数中解析收到的 JSON 指令,控制 LED。
📋 核心结构体一览 (以同步客户端 MQTTClient.h 为例)
我把这些常用的结构体整理了出来,这样理解起来会更清晰:
| 结构体名称 | 主要用途 | 在代码示例中的角色 |
|---|---|---|
MQTTClient |
客户端实例的句柄 (Handle)。它本身是一个不透明的指针,代表了你的 MQTT 客户端对象。 | 在 main() 函数中声明的 MQTTClient client;。后续的 MQTTClient_create() 等函数都通过它来操作这个特定的客户端。 |
MQTTClient_message |
代表一个具体的MQTT消息。它不包含主题(topic),只包含消息的具体内容,如payload(数据负载)、payloadlen(长度)、qos(服务质量)等。 |
出现在 msgarrvd 回调函数中。当设备收到消息时,库会填充好这个结构体并传给你的回调函数。你也用它来构造要发布的消息。 |
MQTTClient_connectOptions |
连接选项。用于配置客户端连接到Broker(MQTT服务器)时的各项参数,如keepAliveInterval(心跳间隔)、cleansession(会话清理)、用户名/密码等。 |
在 main() 函数中声明,并设置好参数后传给 MQTTClient_connect() 发起连接。 |
MQTTClient_willOptions |
遗嘱(Last Will and Testament, LWT)选项。它让你能提前指示Broker,当你的客户端异常断开时,自动代你发布一条预先定义好的消息到某个主题。 | 如果你的设备需要有“上线/离线”通知功能,就会用到这个结构体。它通常作为 MQTTClient_connectOptions 的一个子结构体被使用。 |
MQTTClient_SSLOptions |
TLS/SSL安全连接选项。用于配置如何建立安全的MQTT连接,例如指定CA证书、客户端证书和私钥等。 | 当你连接的Broker要求进行SSL/TLS加密通信时,你需要配置这个结构体,并将其关联到 MQTTClient_connectOptions 中。 |
MQTTClient_createOptions |
客户端创建选项。这是一个更高级的配置结构,主要用于在创建客户端时设定一些不太常用的选项。例如,你可以通过它指定要使用的MQTT协议版本(如5.0)。 | 这是 MQTTClient_create() 函数的高级版本 MQTTClient_createWithOptions() 的参数。 |
MQTTClient_persistence |
消息持久化接口。它不是一个简单的数据容器,而是一个包含函数指针的结构体。它允许你自定义消息持久化的方式,比如将离线期间收到的消息存入文件、数据库或内存中。 | MQTTClient_create() 函数会用到它,你通过传入 MQTTCLIENT_PERSISTENCE_NONE 或 MQTTCLIENT_PERSISTENCE_DEFAULT 来选择默认的持久化行为。 |
MQTTProperties / MQTTProperty |
MQTT 5.0 协议属性。这两个结构体用于支持MQTT 5.0协议中引入的属性(Properties)机制。MQTTProperties 是一个属性的集合。 |
如果你的项目需要使用MQTT 5.0的高级特性(如用户属性、请求-响应模式等),你就需要操作这些结构体。 |
二、逐段详解
1. 头文件和宏定义
c
#include <stdio.h> #include <stdlib.h> #include <string.h> #include <unistd.h> #include "MQTTClient.h" #include "cJSON.h" #include "rpc_client.h" #include "cfg.h" #define ADDRESS "tcp://broker.emqx.io:1883" // 备用 Broker 地址(实际被配置文件覆盖) #define CLIENTID "mqttx_6e973856" // 备用 ClientID(实际被配置文件覆盖) #define TOPIC_SUBSCRIBE "/iot/down" // 备用订阅主题(实际被字符串覆盖) #define TOPIC_PUBLISH "/iot/up" // 备用发布主题(实际被字符串覆盖) #define QOS 1 // 服务质量等级:0/1/2,1表示至少一次送达 #define TIMEOUT 10000L // 等待发布完成的超时时间,单位毫秒 #define USER_NAME "100ask" // 备用用户名(实际被配置文件覆盖) #define PASSWORD "100ask" // 备用密码(实际被配置文件覆盖)
-
MQTTClient.h:Paho MQTT C 客户端库头文件。 -
cJSON.h:轻量级 JSON 解析库。 -
rpc_client.h:自定义 RPC 客户端,用于与底层硬件交互(如rpc_dht11_read、rpc_led_control)。 -
cfg.h:声明read_cfg函数。
注意:代码后面会用 strcpy 将 pub_topic/sub_topic 硬覆盖为这些宏,但实际上动态构造的主题(包含 ProductID/DeviceName)被覆盖掉了,这可能是调试遗留代码。我们分析时仍按功能逻辑理解。
2. 全局变量和回调函数
c
volatile MQTTClient_deliveryToken deliveredtoken;
-
deliveredtoken:标记最近一次消息的发送 token,用于确认消息是否被 Broker 接收。volatile防止编译器优化。token 用于唯一标识某次发布操作,以便在delivered回调中确认是哪一条消息被 Broker 成功接收。delivery为服务器成功接收后,返回给客户端参数的回调函数
2.1 delivered 回调 —— 消息送达确认
void delivered(void *context, MQTTClient_deliveryToken dt)
{
printf("Message with token value %d delivery confirmed\n", dt);
deliveredtoken = dt;
} //confirmed(已确认)
delivered 回调 —— 消息被服务器确认
-
触发时机:当客户端发布的消息(QoS ≥ 1)被服务器成功接收并确认(收到 PUBACK)后,库调用此回调。
-
数据流向:客户端 → 服务器(的确认反馈)。
-
典型用途:确认发布操作完成,可用于日志、释放资源、触发下一批发送。
-
是谁调用的:MQTT 客户端库在收到 PUBACK 后,调用用户注册的
delivered。 -
QoS(Quality of Service,服务质量)0:至多一次 1:至少成功一次 2:只成功一次
-
当发布的消息(QoS≥1)成功被 Broker 确认接收后,库会调用此函数。
-
dt:该消息的唯一 token(由MQTTClient_publishMessage返回)。int 类型 -
作用:让用户知道发送已完成,一般用于日志或触发下一步操作。
-
Message with token value 10 delivery confirmed
如10,为第10次发布
2.2 msgarrvd 回调 —— 收到消息
int msgarrvd(void *context, char *topicName, int topicLen, MQTTClient_message *message)
-
当订阅的主题上有新消息到达时,库自动调用此函数。
-
msgarrvd回调 —— 消息到达客户端 -
触发时机:当客户端从服务器(Broker)收到一条订阅的消息时,MQTT 客户端库会自动调用此回调。 MQTT 客户端库在收到一条订阅的消息时自动调用的函数。
-
数据流向:服务器 → 客户端。
-
典型用途:解析收到的 JSON 指令,执行对应的动作(如控制 LED、上报 ACK)。
-
是谁调用的:MQTT 客户端库在内部收到 PUBLISH 报文后,调用用户注册的
msgarrvd。 -
参数详解:
-
context:用户传入的上下文指针(本例为NULL)。 -
topicName:消息所属的主题字符串。 -
topicLen:主题长度(如果为 0,则需自己strlen)。 -
message:指向消息结构体的指针,包含payload(数据)、payloadlen、qos、retained等。//如发送,用于实际调动开发板上的硬件指令
{
"method": "control",
"clientToken": "mqttx_6e973856",
"params": {
"LED1": 1
}
} //可调动下面代码的LED指令,由客户端发出
-
-
函数内部:
-
打印主题和消息内容。
-
用 cJSON 解析
message->payload。 -
查找
"params"对象,再查找"LED1"。 -
调用
rpc_led_control(1, led1->valueint)控制 LED(valueint为 0 或 1)。 -
释放 cJSON 根对象,以及 MQTT 库分配的消息内存。
-
-
返回值:1 表示成功(非 0 即可)。
关键点:message->payload 是二进制数据,可能不以 '\0' 结尾,但本例中发送的是合法 JSON 字符串,通常发送方会加结束符,所以 cJSON_Parse 直接转换风险较低。更安全的是用 cJSON_ParseWithLength。
2.3 connlost 回调 —— 连接断开
void connlost(void *context, char *cause)
-
当 MQTT 连接意外断开(网络故障、Broker 重启等)时调用。
-
cause:描述断开原因的字符串。 -
本例只打印信息,没有尝试重连(主循环中虽然无限重连,但连接成功后才会进入发布循环。若运行中失去连接,只会打印,程序不会自动重连——这是一个潜在缺陷)。
3. main 函数逐步解析
3.1 变量定义
MQTTClient client; // MQTT 客户端句柄 MQTTClient_connectOptions conn_opts; // 连接选项结构体 int rc; // 返回码,用于检查库函数是否成功 char URI[1000], clientid[1000], username[1000], password[1000], ProductID[1000], DeviceName[1000]; char address[1000]; // 构建的 Broker 地址(形如 "tcp://xxx:1883") char pub_topic[1000], sub_topic[1000]; // 动态构造的发布和订阅主题
3.2 步骤1:读取配置文件
if (0 != read_cfg(URI, clientid, username, password, ProductID, DeviceName))
{
printf("read cfg err\n");
return -1;
}
printf("%s\n%s\n%s\n%s\n%s\n%s\n", URI, clientid, username, password, ProductID, DeviceName);
-
read_cfg解析/etc/wechat.cfg(JSON 格式),将各个字段填入传入的字符数组。 -
如果失败(文件不存在、解析错误等),程序退出。
-
打印出来供调试。
-

3.3 构造连接地址和主题(有逻辑问题)
sprintf(address, "tcp://%s:1883", URI); //假设配置文件中 URI 仅包含 IP 或域名,不含协议和端口 /*此为微信小程序本程序中不需要*/ sprintf(pub_topic, "$thing/up/property/%s/%s", ProductID, DeviceName); sprintf(sub_topic, "$thing/down/property/%s/%s", ProductID, DeviceName); /******************************************************************/ strcpy(pub_topic, TOPIC_PUBLISH); // 硬覆盖为 "/iot/up"发布(主题) strcpy(sub_topic, TOPIC_SUBSCRIBE); // 硬覆盖为 "/iot/down"订阅(主题)
-
正常逻辑:根据配置的 URI(例如
"broker.emqx.io")加上tcp://和端口1883构成完整地址。 -
$thing/up/property/{ProductID}/{DeviceName}是腾讯云 IoT 等平台的标准格式,但随后被strcpy覆盖成简单字符串。最终发布/订阅主题变成了/iot/up和/iot/down。根据实际运行时要根据需求选择使用哪一种。 -
教学解读中,我们忽略覆盖,认为最终使用的是
/iot/up和/iot/down(因为后面MQTTClient_subscribe和MQTTClient_publishMessage用的就是sub_topic、pub_topic)。
3.4 步骤2:初始化 RPC 客户端
if (-1 == RPC_Client_Init())
{
printf("RPC_Client_Init err\n");
return -1;
}
-
连接到本地的 RPC 服务(可能是一个后台进程,提供读写传感器、控制 GPIO 的能力)。
-
失败则退出。
3.5 步骤3:创建 MQTT 客户端
if ((rc = MQTTClient_create(&client, address, clientid, MQTTCLIENT_PERSISTENCE_NONE, NULL)) != MQTTCLIENT_SUCCESS)
-
MQTTClient_create初始化一个客户端实例。-
&client:输出参数,保存句柄。 -
address:Broker 地址字符串(如"tcp://192.168.1.100:1883")。 -
clientid:客户端标识符,必须唯一(Broker 用此区分不同设备)。 -
MQTTCLIENT_PERSISTENCE_NONE:不使用持久化(消息不存磁盘)。 -
NULL:持久化目录(不用)。
-
-
返回值
rc:MQTTCLIENT_SUCCESS(0)表示成功。
3.6 步骤3.1:设置回调函数
MQTTClient_setCallbacks(client, NULL, connlost, msgarrvd, delivered)
-
参数:
-
client:客户端句柄。 -
NULL:上下文指针,会传给各个回调(本例不需要)。 -
connlost:连接丢失回调。 -
msgarrvd:消息到达回调。
-
msgarrvd 回调
你看到的代码中,除了 printf 打印消息内容,还做了:
-
用
cJSON_Parse解析message->payload -
提取
params对象和LED1字段 -
调用
rpc_led_control控制 LED
这些不是库内部实现的,而是(开发者)自己写的业务逻辑。库只负责在收到消息时调用这个回调,并把消息数据通过参数传进来。
delivered:消息已发送回调。
目前只打印了一条确认信息,没有其他动作。你也可以在里面添加自己的逻辑(例如记录日志、释放资源、触发下一次发送等)。
-
必须在连接之前设置,回调才能生效。
3.7 步骤3.2:配置连接参数
conn_opts.keepAliveInterval = 100; // 保活周期,单位秒。100秒内无任何报文则发送 PINGREQ conn_opts.cleansession = 1; // 1 = 清除之前的会话(订阅等),0 = 恢复之前会话 conn_opts.username = username; // 用户名指针 conn_opts.password = password; // 密码指针
-
conn_opts必须先初始化为MQTTClient_connectOptions_initializer,再修改所需字段。 -
如果 Broker 不需要认证,可以设为
NULL。
3.8 步骤4:连接服务器(带重试)
while (1)
{
if ((rc = MQTTClient_connect(client, &conn_opts)) != MQTTCLIENT_SUCCESS)
{
printf("Failed to connect, return code %d\n", rc);
// 不退出,继续重试
sleep(2);
}
else
{
break;
}
}
-
MQTTClient_connect发起 TCP 连接和 MQTT 握手。 -
失败则打印并循环(没有延时,会疯狂重试,最好加
sleep(1))。 -
成功后跳出循环。
3.9 步骤5:订阅主题
c
if ((rc = MQTTClient_subscribe(client, sub_topic, QOS)) != MQTTCLIENT_SUCCESS)
{
printf("Failed to subscribe, return code %d\n", rc);
rc = EXIT_FAILURE;
}
else
{
// 进入发布循环
}
-
MQTTClient_subscribe发送 SUBSCRIBE 报文。 -
sub_topic:要订阅的主题(如/iot/down)。 -
QOS:订阅的 QoS 等级(本例为 1,Broker 会保证至少一次送达)。 -
订阅成功后,后续发布到该主题的消息会触发
msgarrvd。
逻辑原因
1. 确保能收到下行消息
-
设备需要接收云端或手机端发来的控制命令(如开关 LED)。
-
这些命令可能随时到达,设备必须在连接成功后尽快注册自己的兴趣(即订阅主题),否则一旦有命令在设备订阅之前到达,设备就会错过。
-
先订阅可以避免消息丢失(尤其是 QoS ≥ 1 时,服务器会为离线客户端保留消息的前提是会话被持久化,但你的代码中
cleansession=1,会话是全新的,先订阅才能保证后续消息到达)。
2. 发布是周期性的,订阅是一次性的
-
发布(上报传感器数据)是在主循环中每隔 5 秒执行一次,可以随时开始。
-
订阅通常只需要在连接成功后执行一次,之后一直有效。
-
先完成订阅,再进入发布循环,符合逻辑顺序。
3. 大多数物联网场景的典型模式
-
设备上电 → 连接 Broker → 订阅下行主题 → 进入主循环(发布上行数据 + 处理下行命令)。
-
如果不先订阅,可能在发布几次之后才能收到命令,造成响应延迟。
技术上的可行性
-
MQTT 协议本身不强制订阅必须在发布之前。你可以先发布后订阅,甚至只发布不订阅。
-
但如果设备需要接收消息,就必须在对方发送消息之前完成订阅。由于你无法预测控制命令何时下发,最保险的做法就是连接成功后立即订阅。
3.10 发布循环内的操作
char humi, temp; while (0 != rpc_dht11_read(&humi, &temp)); // 阻塞直到读取成功
-
rpc_dht11_read通过 RPC 读取温湿度,返回 0 表示成功。 -
humi和temp是char类型,湿度一般为 0~100,温度可能负值?实际应使用int,但这里示例用了char,可能有截断风险。
然后构造 JSON:
c
sprintf(buf, "{\"method\":\"report\",\"clientToken\":\"mqttx_6e973856\",\"timestamp\":1628646783,\"params\":{\"temp_value\":%d,\"humi_value\":%d}}", temp, humi);
-
使用
sprintf将 JSON 写入buf。注意时间戳是固定的,实际应该动态获取。
填充发布消息结构体:
c
pubmsg.payload = buf; // 指向 JSON 字符串 pubmsg.payloadlen = (int)strlen(buf); pubmsg.qos = QOS; // 1 pubmsg.retained = 0; // 不保留消息,即不保存最新值给后续订阅者
发布:
if ((rc = MQTTClient_publishMessage(client, pub_topic, &pubmsg, &token)) != MQTTCLIENT_SUCCESS)
{
printf("Failed to publish message, return code %d\n", rc);
continue;
}
-
MQTTClient_publishMessage:将消息发送到 Broker。 -
pub_topic:主题(/iot/up)。 -
&token:输出参数,代表本次发布的传输令牌。 -
返回值非 0 表示发布请求失败(网络断开、无效参数等)。
等待发布完成(仅对 QoS≥1 有意义):
c
rc = MQTTClient_waitForCompletion(client, token, TIMEOUT);
printf("Message with delivery token %d delivered\n", token);
-
MQTTClient_waitForCompletion阻塞直到该 token 对应的消息被确认(收到 PUBACK)或超时。 -
TIMEOUT为 10000 毫秒。 -
对于 QoS 0,该函数会立即返回成功。
最后 sleep(5),每 5 秒上报一次。
3.11 清理(正常不会执行到)
c
MQTTClient_unsubscribe(client, sub_topic); MQTTClient_disconnect(client, 10000); MQTTClient_destroy(&client);
-
由于发布循环是
while(1),这些代码不会执行。仅当订阅失败或循环被 break 时才会走到。实际产品中应加入退出机制(如检测到'Q'按键)。
三、参数配合和整体框架
各组件如何配合?
-
MQTTClient_create+MQTTClient_setCallbacks
→ 创建客户端并注册回调,为后续收发消息做准备。 -
MQTTClient_connect
→ 建立网络连接和 MQTT 会话,Broker 会验证clientid/username/password。 -
MQTTClient_subscribe
→ 告诉 Broker:“请把主题sub_topic的消息转发给我”。
之后 Broker 收到匹配sub_topic的消息时,就会推送给该客户端。 -
msgarrvd回调
→ 当 Broker 推送消息到达时,库自动调用它。在回调内解析 JSON 并执行动作(控制 LED)。 -
MQTTClient_publishMessage
→ 发布消息到pub_topic,Broker 会转发给所有订阅了该主题的客户端。 -
MQTTClient_waitForCompletion
→ 确保(QoS≥1)消息已正确送达,否则可重试。
回调函数与主循环的关系
-
主循环负责主动发布传感器数据。
-
回调函数
msgarrvd由库内部线程(或主循环中的MQTTClient_yield,本例未显式调用)触发。Paho 库通常会在MQTTClient_connect和publishMessage等函数内部自行处理网络 I/O 并分发回调。所以不需要额外的MQTTClient_receive调用。
四、整体框架总结(可复用的模板)
以下是编写类似 MQTT 客户端程序的通用步骤:
text
1. 读取配置(Broker地址、ClientID、用户名密码、业务主题参数) 2. 初始化硬件接口(RPC、GPIO、传感器驱动) 3. 创建 MQTT 客户端(MQTTClient_create) 4. 设置回调(connlost, msgarrvd, delivered) 5. 配置连接选项(keepAlive, cleansession, 用户名密码) cleansession会话状态的处理方式
| 标志值 | 会话行为 | 订阅是否持久 | 离线期间的消息是否接收 |
|---|---|---|---|
cleansession = 1 |
每次连接都是全新会话,之前的订阅、未确认的消息等全部丢弃 | 每次连接后必须重新订阅 | 不接收离线期间的消息 |
cleansession = 0 |
Broker 记住会话,包括订阅、未完成的 QoS 1/2 消息 | 订阅持续有效,断线重连后无需重新订阅 | 设备离线时到达的消息,Broker 会保留(取决于 QoS),重连后会收到 |
6. 循环连接直到成功(MQTTClient_connect) 7. 订阅下行主题(MQTTClient_subscribe) 8. 进入主循环: - 读取传感器数据 - 构造 JSON 消息 - 发布到上行主题(MQTTClient_publishMessage) - 等待送达确认(可选) - 延时或等待下次周期 9. (如果退出)取消订阅、断开连接、销毁客户端
五、编写类似程序时需要注意的要点
-
缓冲区大小:用宏定义统一长度,避免硬编码 1000。
-
JSON 解析:使用
cJSON_ParseWithLength传入payloadlen,防止缺'\0'。 -
连接重试:重试循环内应加入延迟(
sleep(1)),避免 CPU 彪高。 -
断开重连机制:在
connlost回调中设置标志,主循环检测后重新调用MQTTClient_connect。 -
资源释放:在
msgarrvd中必须调用MQTTClient_freeMessage和MQTTClient_freeTopicName防止内存泄漏。 -
主题设计:区分上行(设备→云)和下行(云→设备),保持一致性。
-
线程安全:如果主线程和多线程回调共享数据,需加锁。本例只有单线程回调(Paho 默认在调用
MQTTClient_connect/publish的线程中处理 I/O),相对安全。
希望这个详尽的分析能帮助你完全理解代码,并自己编写出健壮的 MQTT 设备端程序。
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include "MQTTClient.h"
#include "cJSON.h"
#include "rpc_client.h"
#include "cfg.h"
#define ADDRESS "tcp://broker.emqx.io:1883"
#define CLIENTID "mqttx_6e973856"
#define TOPIC_SUBSCRIBE "/iot/down"
#define TOPIC_PUBLISH "/iot/up"
#define QOS 1
#define TIMEOUT 10000L
#define USER_NAME "100ask"
#define PASSWORD "100ask"
volatile MQTTClient_deliveryToken deliveredtoken;
//送达至服务器成功
void delivered(void *context, MQTTClient_deliveryToken dt)
{
printf("Message with token value %d delivery confirmed\n", dt);
deliveredtoken = dt;
}
//传给客户端成功 //由 MQTT 客户端库在收到消息时自动调用
int msgarrvd(void *context, char *topicName, int topicLen, MQTTClient_message *message)
{
printf("Message arrived\n");
printf(" topic: %s\n", topicName);
printf(" message: %.*s\n", message->payloadlen, (char*)message->payload);
/* 根据接收到的消息控制LED */
/* 消息格式: {"method":"control","clientToken":"v2177557380FdRVD::DU2S%vwG7","params":{"LED1":0}} */
cJSON *root = cJSON_Parse((char*)message->payload); //有效数据
cJSON *ptTemp = cJSON_GetObjectItem(root, "params");
if (ptTemp)
{
cJSON *led1 = cJSON_GetObjectItem(ptTemp, "LED1");
if (led1)
{
rpc_led_control(1,led1->valueint);
}
}
cJSON_Delete(root);
MQTTClient_freeMessage(&message);
MQTTClient_free(topicName);
return 1;
}
/*{
* "method": "control",
* "clientToken": "mqttx_6e973856",
* "params": {
* "LED1": 1
* }
*}
*/ //下行设备(客户端)手机,可发送此命令控制LED
void connlost(void *context, char *cause)
{
printf("\nConnection lost\n");
if (cause)
printf(" cause: %s\n", cause);
}
int main(int argc, char* argv[])
{
MQTTClient client;
MQTTClient_connectOptions conn_opts = MQTTClient_connectOptions_initializer;
int rc;
char URI[1000];
char clientid[1000];
char username[1000];
char password[1000];
char ProductID[1000];
char DeviceName[1000];
char address[1000];
//发布主题+接收主题
char pub_topic[1000]; /* $thing/up/property/{ProductID}/{DeviceName} */
char sub_topic[1000]; /* $thing/down/property/{ProductID}/{DeviceName} */
/* 1. 读配置文件 /etc/wechat.cfg得到URI,CLIENTID,USERNAME,PASSWD,ProductID,DeviceName */
if ( 0 != read_cfg(URI, clientid, username, password, ProductID, DeviceName))
{
printf("read cfg err\n");
return -1;
}
printf("%s\n%s\n%s\n%s\n%s\n%s\n",URI,clientid, username, password, ProductID, DeviceName);
//strcpy(URI, ADDRESS);
//strcpy(clientid, CLIENTID);
//strcpy(username, "100ask");
//strcpy(password, "100ask");
//sprintf(address, "%s", URI);
sprintf(address, "tcp://%s:1883", URI); //产品ID + 设备名称
//属性
//sprintf(pub_topic, "$thing/up/property/%s/%s", ProductID, DeviceName);
//sprintf(sub_topic, "$thing/down/property/%s/%s", ProductID, DeviceName);
strcpy(pub_topic, TOPIC_PUBLISH);
strcpy(sub_topic, TOPIC_SUBSCRIBE);
/* INIT RPC: connect RPC Server */
if (-1 == RPC_Client_Init()) // 2. 启动RPC客户端
{
printf("RPC_Client_Init err\n");
return -1;
}
// 3. 创建MQTT客户端
if ((rc = MQTTClient_create(&client, address, clientid,MQTTCLIENT_PERSISTENCE_NONE, NULL)) != MQTTCLIENT_SUCCESS)
{
printf("Failed to create client, return code %d\n", rc);
rc = EXIT_FAILURE;
goto exit;
}
// 3.1设置回调函数,打印客户端在连接过程中的错误
if ((rc = MQTTClient_setCallbacks(client, NULL, connlost, msgarrvd, delivered)) != MQTTCLIENT_SUCCESS)
{
printf("Failed to set callbacks, return code %d\n", rc);
rc = EXIT_FAILURE;
goto destroy_exit;
}
// 3.2配置客户端连接参数
conn_opts.keepAliveInterval = 100;
conn_opts.cleansession = 1;
conn_opts.username = username;
conn_opts.password = password;
while (1)
{
// 4.客户端连接服务器
if ((rc = MQTTClient_connect(client, &conn_opts)) != MQTTCLIENT_SUCCESS)
{
printf("Failed to connect, return code %d\n", rc);
rc = EXIT_FAILURE;
sleep(1);
//goto destroy_exit;
}
else
{
break;
}
}
printf("Subscribing to topic %s\nfor client %s using QoS%d\n\n"
"Press Q<Enter> to quit\n\n", TOPIC_SUBSCRIBE, CLIENTID, QOS);
//5.1订阅消息
if ((rc = MQTTClient_subscribe(client, sub_topic, QOS)) != MQTTCLIENT_SUCCESS)
{
printf("Failed to subscribe, return code %d\n", rc);
rc = EXIT_FAILURE;
}
else
{
int ch;
int cnt = 0;
//初始化订阅服务
MQTTClient_message pubmsg = MQTTClient_message_initializer;
char buf[1000];
MQTTClient_deliveryToken token; //送达
while (1)
{
/* rpc_dht11_read */
char humi, temp;
while (0 !=rpc_dht11_read(&humi, &temp));
sprintf(buf, "\
{ \
\"method\":\"report\",\
\"clientToken\":\"mqttx_6e973856\",\
\"timestamp\":1628646783,\
\"params\":{\
\"temp_value\":%d,\
\"humi_value\":%d\
}\
}", temp, humi);
pubmsg.payload = buf;
pubmsg.payloadlen = (int)strlen(buf);
pubmsg.qos = QOS;
pubmsg.retained = 0;
/* 发布消息 */
//6.发布消息至服务器
if ((rc = MQTTClient_publishMessage(client, pub_topic, &pubmsg, &token)) != MQTTCLIENT_SUCCESS)
{
printf("Failed to publish message, return code %d\n", rc);
continue;
}
//等待发布完成
rc = MQTTClient_waitForCompletion(client, token, TIMEOUT);
printf("Message with delivery token %d delivered\n", token);
sleep(5);
}
//取消订阅
if ((rc = MQTTClient_unsubscribe(client, sub_topic)) != MQTTCLIENT_SUCCESS)
{
printf("Failed to unsubscribe, return code %d\n", rc);
rc = EXIT_FAILURE;
}
}
if ((rc = MQTTClient_disconnect(client, 10000)) != MQTTCLIENT_SUCCESS)
{
printf("Failed to disconnect, return code %d\n", rc);
rc = EXIT_FAILURE;
}
destroy_exit:
MQTTClient_destroy(&client);
exit:
return rc;
}

#include <sys/types.h>
#include <sys/stat.h>
#include <fcntl.h>
#include <stdio.h>
#include <unistd.h>
#include "cfg.h"
#include "cJSON.h"
#define CFG_STR_LEN 256
int read_cfg(char *URI,char *clientid,char *username,char *password,char *ProductID,char *DeviceName)
{
int ret;
cJSON *ptTemp;
char buf[1024];
int fd = open(CFG_FILE,O_RDONLY);
if(fd < 0)
{
perror("open(CFG_FILE) ");
return -1;
}
ret = read(fd,buf,sizeof(buf));
if(ret < 0)
{
close(fd);
return -1;
}
else
buf[ret] = '\0';
cJSON *root = cJSON_Parse(buf);
ptTemp = cJSON_GetObjectItem(root,"URI");
if(ptTemp)
snprintf(URI, CFG_STR_LEN,"%s",ptTemp->valuestring);
ptTemp = cJSON_GetObjectItem(root,"clientid");
if(ptTemp)
snprintf(clientid, CFG_STR_LEN,"%s",ptTemp->valuestring);
ptTemp = cJSON_GetObjectItem(root,"username");
if(ptTemp)
snprintf(username, CFG_STR_LEN,"%s",ptTemp->valuestring);
ptTemp = cJSON_GetObjectItem(root,"password");
if(ptTemp)
snprintf(password, CFG_STR_LEN,"%s",ptTemp->valuestring);
ptTemp = cJSON_GetObjectItem(root,"ProductID");
if(ptTemp)
snprintf(ProductID, CFG_STR_LEN,"%s",ptTemp->valuestring);
ptTemp = cJSON_GetObjectItem(root,"DeviceName");
if(ptTemp)
snprintf(DeviceName, CFG_STR_LEN,"%s",ptTemp->valuestring);
return 0;
}

更多推荐
所有评论(0)