Paho MQTTAsync_sendMessage数据传递与内存管理方案咨询
实现方案
你不需要用全局变量,基础场景下也不需要额外维护待发数据列表,核心思路是利用Paho MQTT异步接口所有回调支持传入自定义context指针的特性,封装独立的发送上下文结构体管理每一条消息的资源,生命周期和单次发送请求绑定即可。
核心实现逻辑
- 定义自定义发送上下文结构体,封装MQTT客户端实例、拷贝后的payload、payload长度等资源,每一次
send()调用生成一个独立的上下文实例 - 在
send()函数中拷贝传入的data数据到上下文的payload字段,避免原指针释放导致野指针 - 所有回调(连接成功/失败、发送成功/失败)都传入该上下文实例,在最终回调(发送结束/连接失败)中统一释放上下文和payload的内存,实现自动内存管理
完整代码示例
1. 自定义上下文结构与公共定义
#include <stdio.h> #include <stdlib.h> #include <string.h> #include "MQTTAsync.h" #define QOS 1 #define DEFAULT_TOPIC "your/default/topic" // 自定义发送上下文,每一条待发消息对应一个独立实例 typedef struct { MQTTAsync client; char* payload; size_t payload_len; const char* topic; // 可按需扩展支持动态topic } SendContext; // 提前创建并初始化好的全局MQTT客户端实例,避免每次发送新建连接 static MQTTAsync g_mqtt_client = NULL;
2. 发送逻辑封装与回调实现
// 发送完成(成功/失败)统一释放资源 static void release_send_context(SendContext* ctx) { if (ctx == NULL) return; if (ctx->payload != NULL) { free(ctx->payload); } free(ctx); } // 实际发送消息逻辑,连接正常时直接调用,连接成功后回调调用 static void do_send(SendContext* ctx) { MQTTAsync_responseOptions opts = MQTTAsync_responseOptions_initializer; MQTTAsync_message pubmsg = MQTTAsync_message_initializer; int rc; opts.onSuccess = on_send_success; opts.onFailure = on_send_failure; opts.context = ctx; pubmsg.payload = ctx->payload; pubmsg.payloadlen = ctx->payload_len; pubmsg.qos = QOS; pubmsg.retained = 0; if ((rc = MQTTAsync_sendMessage(ctx->client, ctx->topic, &pubmsg, &opts)) != MQTTASYNC_SUCCESS) { printf("Failed to start sendMessage, return code %d\n", rc); release_send_context(ctx); } } // 连接成功回调 void on_connect_success(void* context, MQTTAsync_successData* response) { printf("Connection success\n"); SendContext* ctx = (SendContext*)context; do_send(ctx); } // 连接失败回调 void on_connect_failure(void* context, MQTTAsync_failureData* response) { printf("Connect failed, rc %d\n", response ? response->code : -1); SendContext* ctx = (SendContext*)context; release_send_context(ctx); } // 发送成功回调 void on_send_success(void* context, MQTTAsync_successData* response) { printf("Message send success\n"); release_send_context((SendContext*)context); } // 发送失败回调 void on_send_failure(void* context, MQTTAsync_failureData* response) { printf("Message send failed, rc %d\n", response ? response->code : -1); release_send_context((SendContext*)context); }
3. 对外暴露的send函数实现
void send(const char* data) { if (data == NULL || g_mqtt_client == NULL) return; // 分配上下文内存 SendContext* ctx = (SendContext*)malloc(sizeof(SendContext)); if (ctx == NULL) { printf("Malloc SendContext failed\n"); return; } // 拷贝payload,和原输入指针解耦 ctx->payload_len = strlen(data); ctx->payload = (char*)malloc(ctx->payload_len + 1); if (ctx->payload == NULL) { printf("Malloc payload failed\n"); free(ctx); return; } strcpy(ctx->payload, data); ctx->client = g_mqtt_client; ctx->topic = DEFAULT_TOPIC; // 检查连接状态 if (MQTTAsync_isConnected(g_mqtt_client)) { // 已连接直接发送 do_send(ctx); } else { // 未连接先发起连接 MQTTAsync_connectOptions conn_opts = MQTTAsync_connectOptions_initializer; conn_opts.keepAliveInterval = 20; conn_opts.cleansession = 1; conn_opts.onSuccess = on_connect_success; conn_opts.onFailure = on_connect_failure; conn_opts.context = ctx; int rc; if ((rc = MQTTAsync_connect(g_mqtt_client, &conn_opts)) != MQTTASYNC_SUCCESS) { printf("Failed to start connect, return code %d\n", rc); release_send_context(ctx); } } }
待发列表的适用场景
只有当你需要支持断网缓存补发的能力时,才需要实现线程安全的待发列表:调用send()时先把消息存入列表,网络恢复后按顺序取出列表中的消息发送,发送成功再从列表中移除并释放资源。如果不需要断网缓存能力,上述实现已经完全满足你需要的即发即忘、自动内存管理的需求。
注意事项
- 示例中的
g_mqtt_client需要提前调用MQTTAsync_create完成初始化,配置好服务器地址、客户端ID等参数 - 如果多线程调用
send()函数,需要对MQTT客户端的连接、发送操作加锁,避免并发冲突 - 如果需要支持动态传入topic、qos参数,可以扩展
send函数的入参,将参数存入SendContext结构体即可
内容的提问来源于stack exchange,提问作者vico
相关产品推荐
相关产品推荐

