You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Mosquitto的MQTT客户端Wrapper消息发送阻塞问题求助

问题根源与修复方案

核心问题分析

  1. 回调线程阻塞死锁:on_message回调运行在Mosquitto的Loop线程中,你在回调里调用client->Publish(),而Publish里的while(handshake.load() != ACK_BYTE_RECEIVE)是死循环阻塞,直接卡住了Loop线程——Mosquitto无法再处理任何新消息(包括对方的ACK回复),导致握手永远无法完成,后续消息全部失败。
  2. 全局状态变量污染:toSelf是全局变量,多订阅主题或并发消息时,状态会被覆盖,导致握手逻辑混乱,多订阅场景下直接触发程序挂起。
  3. 冗余自定义握手: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.23 09:14:53