基于Mosquitto的MQTT客户端Wrapper消息发送阻塞问题求助
问题根源与修复方案
核心问题分析
- 回调线程阻塞死锁:
on_message回调运行在Mosquitto的Loop线程中,你在回调里调用client->Publish(),而Publish里的while(handshake.load() != ACK_BYTE_RECEIVE)是死循环阻塞,直接卡住了Loop线程——Mosquitto无法再处理任何新消息(包括对方的ACK回复),导致握手永远无法完成,后续消息全部失败。 - 全局状态变量污染:
toSelf是全局变量,多订阅主题或并发消息时,状态会被覆盖,导致握手逻辑混乱,多订阅场景下直接触发程序挂起。 - 冗余自定义握手:MQTT本身的QoS机制已能保证消息可靠交付,自定义ACK握手不仅多余,还引入了线程安全风险。
具体修复步骤
1. 替换Publish中的死循环阻塞,改用异步等待
如果必须保留自定义握手,用条件变量替代死循环,避免阻塞消息处理线程:
// Client类中添加成员变量 #include <condition_variable> #include <mutex> class Client { private: std::condition_variable cv; std::mutex mtx; bool toSelf = false; // 把全局toSelf移为类成员 // ...其他已有成员 }; // 修改Publish方法 int Client::Publish(Message msg, const char* topic) { std::unique_lock<std::mutex> lock(mtx); handshake.store(ACK_BYTE_SEND); AckSend(topic); // 带超时的等待,避免永久阻塞 if(cv.wait_for(lock, std::chrono::seconds(5), [this]{ return handshake.load() == ACK_BYTE_RECEIVE; })) { handshake.store(-1); // 执行消息发布 char serialized_buf[sizeof(Message)]; std::memcpy(serialized_buf, &msg, sizeof(Message)); printf("[%s] Published message to topic %s with commando \"%d\"\n", this->client_id, topic, msg.commando); return mosquitto_publish(mosq, NULL, topic, sizeof(serialized_buf), (void*)serialized_buf, 2, false); } else { printf("Handshake timeout for topic %s\n", topic); handshake.store(-1); return -1; } } // 修改on_message中的ACK_RECEIVE分支,唤醒等待线程 else if (payload == ACK_BYTE_RECEIVE) { printf("test3\n"); client->handshake.store(ACK_BYTE_RECEIVE); client->cv.notify_one(); // 唤醒Publish的等待逻辑 return; }
2. 修复回调中传递Client实例的方式
避免使用全局client指针,通过Mosquitto的userdata传递Client实例,确保线程安全:
// 修改Connect方法,将当前Client实例作为userdata传入 mosq = mosquitto_new(client_id, true, this); // 替换原NULL为this // 修改on_message回调,从userdata获取Client实例 void on_message(struct mosquitto* mosq, void *obj, const struct mosquitto_message *msg) { Client* client = static_cast<Client*>(obj); if (mosq == nullptr) { printf("Variable mosq is a nullptr\n"); exit(-1); } // BEGIN HANDSHAKE // int payload = *(int*)msg->payload; if (payload == ACK_BYTE_SEND && client->handshake.load() == ACK_BYTE_SEND) { printf("test1\n"); client->toSelf = true; return; } else if (client->handshake.load() == -1) { printf("test2\n"); if (payload == ACK_BYTE_SEND) { client->AckReceive(msg->topic); return; } else if (payload == ACK_BYTE_RECEIVE) { return; } } else if (payload == ACK_BYTE_RECEIVE) { printf("test3\n"); client->handshake.store(ACK_BYTE_RECEIVE); client->cv.notify_one(); return; } if (client->toSelf) { printf("test4\n"); client->toSelf = false; client->handshake.store(-1); return; } printf("test5\n"); // END HANDSHAKE // // BEGIN USER CODE // struct Message bericht; std::memcpy(&bericht, msg->payload, msg->payloadlen); if (!strcmp(msg->topic, "hmi/warehouse")) { switch (bericht.commando) { case HMI_TO_WAREHOUSE__MESSAGE_DATA: printf("[%s] received HMI_TO_WAREHOUSE__MESSAGE_DATA.\n", client->GetClientID()); bericht.commando = WAREHOUSE_TO_CRANE__READY_FOR_PICKUP; // 启动独立线程执行Publish,避免阻塞回调线程 std::thread([client, bericht, topic=std::string("warehouse/crane")](){ client->Publish(bericht, topic.c_str()); }).detach(); break; } } // END USER CODE }
3. 多订阅场景下的状态隔离
如果需要同时处理多个主题的握手,不要用全局的handshake变量,改为给每个主题维护独立的状态(比如用std::unordered_map<std::string, std::atomic<int>> topic_handshakes),避免不同主题的握手状态互相干扰。
额外建议
- 优先使用MQTT原生QoS机制(你已设置QoS 2),删除自定义握手逻辑,减少不必要的复杂度。
- 在
Loop方法中处理mosquitto_loop的返回值,出现错误时及时重连或处理,避免消息循环中断:
void Client::Loop() { while (!ShouldExit()) { int res = mosquitto_loop(mosq, 1000, 60); if (res != MOSQ_ERR_SUCCESS && res != MOSQ_ERR_NO_CONN) { printf("Mosquitto loop error: %d\n", res); // 这里可以添加重连逻辑 } } }
内容的提问来源于stack exchange,提问作者user25291833
相关产品推荐
相关产品推荐

