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

使用AMQP-CPP与Boost Asio时随机出现Connection lost错误排查

问题

使用AMQP-CPP库(RabbitMQ的C++封装)搭建RPC服务端及多客户端,采用Boost Asio作为事件循环并运行在std::thread中。实现了RabbitMQ封装类负责队列、交换机、服务端、客户端等操作。在测试用例中重复1000次创建RPC客户端并连接服务端的操作,随机出现**"Channel open error - connection lost"**错误(错误出现的时机和迭代次数不固定,可能在队列声明后、收到响应后等任意步骤)。

测试代码

RabbitMQ srv_rabbitmq(uri);

srv_rabbitmq.start();
// 声明队列、交换机并绑定
auto rpc_server = srv_rabbitmq.make_rpc_server(server_config, rpc_server_cb);

auto rabbit_f = [&client_config, &request]() {
    try {
        RabbitMQ rabbitmq(uri);
        rabbitmq.start();
        auto rpc_client = rabbitmq.make_rpc_client(client_config);
        auto response = rpc_client->send_request(request);
        ASSERT_EQ(response.Body, request.Body);
    } TEST_CATCH
};
for (auto i = 0; i < test_count; ++i) {
    rabbit_f();        
    std::cout << "Successfully connected " << (i + 1) << "/" << test_count << std::endl;
}

RabbitMQ类核心代码

start()方法

void RabbitMQ::start() {
    {
        // 其他代码
        m_handler = std::make_unique<AMQP::LibBoostAsioHandler>(service);
        m_connection = std::make_shared<AMQP::TcpConnection>(m_handler.get(), AMQP::Address(m_uri));
    }    
    m_working_thread = std::thread([&] {        
        // 其他代码
        while(service.run());    
    });
}

析构函数与stop()方法

RabbitMQ::~RabbitMQ() {
    if (m_working_thread.joinable()) {
        stop();        
        m_working_thread.join();
    }    
    if (!service.stopped()) {
        m_handler.reset();
    }
}

void RabbitMQ::stop() {
    if (m_is_run) {
        service.stop();
        m_connection->close();
        service.reset();
    }
}

测试日志

Successfully connected 65/1000
RabbitMQ RPC Client - success declare queue
[2024-04-18 17:05:51.877] [error] "caller":"/home/vlad/Maynitek/rabbitmq-lib/src/RabbitMQ/RPC/Client.cpp:137","error":"connection lost","exchange":"e.test.rpc","func":"operator()","msg":"RabbitMQ RPC Client - error","queue":"\u0010>\u0000XJw\u0000\u0000\u000b\u0000\u0000\u0000\u0000\u0000\u0000\u0000routing_key\u0000Hw","routing_key":"","rpc_namespace":"mongoose::test"

问题原因与修复方案

核心问题点

  1. 线程捕获与io_context生命周期冲突:start()中线程lambda用[&]捕获service,且stop()中先调用service.stop()再直接service.reset(),此时线程还在执行while(service.run()),会访问已被重置的io_context,触发未定义行为,导致连接异常断开。
  2. 连接关闭与事件循环顺序颠倒:stop()先停止io_context再关闭AMQP连接,此时handler可能还在处理异步操作,强行停止会导致连接未正常释放,残留资源干扰后续新连接。
  3. handler释放时机错误:析构函数在service未完全停止时释放m_handler,此时handler可能仍被io_context的异步任务引用,引发内存访问错误。

修复步骤

1. 修正线程捕获方式

改用this捕获成员变量,避免悬空引用:

m_working_thread = std::thread([this] {        
    // 其他代码
    while (!service.stopped() && service.run());    
});

2. 调整stop()方法执行顺序

先关闭连接,再停止io_context,最后等待线程退出后重置:

void RabbitMQ::stop() {
    if (!m_is_run) return;
    m_is_run = false;

    // 优先关闭AMQP连接,让handler完成清理
    if (m_connection) {
        m_connection->close();
    }

    // 停止io_context,不再接受新任务
    service.stop();

    // 等待线程完全退出
    if (m_working_thread.joinable()) {
        m_working_thread.join();
    }

    service.reset();
}

3. 优化析构函数逻辑

确保handler在所有资源释放后再销毁:

RabbitMQ::~RabbitMQ() {
    stop();
    // 此时service已停止、线程已退出,安全释放handler
    m_handler.reset();
}

4. 客户端同步逻辑补充

在send_request完成后,确保所有异步操作执行完毕再销毁RabbitMQ对象,避免未完成的请求导致连接异常。

额外建议

  • 高频创建销毁客户端场景,建议使用连接池,减少资源开销与连接抖动。
  • 检查RabbitMQ服务器的最大连接数配置,避免短时间内创建过多连接触发阈值。
  • 在AMQP-CPP的handler中添加更详细的错误回调,精准定位连接断开的具体原因。

内容的提问来源于stack exchange,提问作者Vladimir Privalov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 20:43:26