使用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"
问题原因与修复方案
核心问题点
- 线程捕获与io_context生命周期冲突:
start()中线程lambda用[&]捕获service,且stop()中先调用service.stop()再直接service.reset(),此时线程还在执行while(service.run()),会访问已被重置的io_context,触发未定义行为,导致连接异常断开。 - 连接关闭与事件循环顺序颠倒:
stop()先停止io_context再关闭AMQP连接,此时handler可能还在处理异步操作,强行停止会导致连接未正常释放,残留资源干扰后续新连接。 - 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
相关产品推荐
相关产品推荐

