基于Boost.Asio与Boost.Interprocess Message Queue的服务器异步读写阻塞问题求解
解决Boost Asio异步读阻塞导致无法处理消息队列的问题
你的问题核心在于Asio的事件循环被TCP异步读操作占据,同时没有持续监听消息队列的新数据,导致服务器没法及时把消息队列里的响应写回TCP连接。我来帮你拆解问题并给出可行的解决方案。
问题分析
先看你的代码里几个关键的问题点:
- do_write只执行一次:在
start()里调用了一次do_write(),但如果消息队列一开始没有数据,try_receive直接返回,之后就再也不会主动检查消息队列了。即使后面消息队列有新的响应数据,也没有触发写操作的逻辑。 - 异步读的回调逻辑有缺陷:
do_read_header()的回调里不管有没有错误(比如socket断开),都无条件调用do_read_header(),这会导致无效循环,甚至占用strand的执行时间,影响其他操作。 - 消息队列与Asio事件循环脱节:Boost进程间消息队列的操作不属于Asio的IO事件,Asio的io_context不会主动感知到消息队列的新数据,必须我们主动轮询或用线程监听。
可行解决方案
这里提供两种常用思路,你可以根据业务场景选择:
方案1:用Asio定时器定时轮询消息队列
这种方案不需要额外线程,通过定时器定期检查消息队列,适合消息队列数据不是特别频繁的场景。
修改你的Session类,关键改动如下:
class Session : public std::enable_shared_from_this<Session> { public: Session(tcp::socket socket) : socket_(std::move(socket)) , response_queue_(boost::interprocess::open_or_create, "response_queue", 100, 100) , request_queue_(boost::interprocess::open_or_create, "request_queue", 100, 100) , timer_(socket_.get_executor()) // 复用socket的执行器初始化定时器 { } void start() { post(strand_, [this, self = shared_from_this()] { start_response_poll(); // 启动消息队列轮询 do_read_header(); }); } private: void start_response_poll() { auto self(shared_from_this()); // 每隔100ms轮询一次(时间间隔可根据需求调整) timer_.expires_after(std::chrono::milliseconds(100)); timer_.async_wait(strand_.wrap([this, self](boost::system::error_code ec) { if (!ec && socket_.is_open()) { do_write(); start_response_poll(); // 继续下一轮轮询 } })); } void do_write() { auto self(shared_from_this()); Message msg; boost::interprocess::message_queue::size_type recvd_size; unsigned int priority = 0; // 有数据才处理写操作 if(response_queue_.try_receive(msg.body(), msg.max_body_length, recvd_size, priority)) { try { msg.body_length(recvd_size); msg.encode_header(); std::memcpy(data_, msg.data(), msg.length()); } catch (std::exception& e ) { std::cout << e.what(); return; } std::cout << "write " << msg.body() << "\n"; boost::asio::async_write(socket_, boost::asio::buffer(data_, msg.length()), strand_.wrap([this, self](boost::system::error_code ec, std::size_t /*length*/) { if (ec) { std::cerr << "write error:" << ec.value() << " message: " << ec.message() << "\n"; socket_.close(); timer_.cancel(); // 关闭定时器 } })); } } // 修复异步读头部的回调逻辑 void do_read_header() { auto self(shared_from_this()); std::cout << "do_read_header\n"; boost::asio::async_read(socket_, boost::asio::buffer(res.data(), res.header_length), strand_.wrap([this, self](boost::system::error_code ec, std::size_t /*length*/) { if (ec) { std::cerr << "read header error:" << ec.value() << " message: " << ec.message() << "\n"; socket_.close(); timer_.cancel(); return; } if (res.decode_header()) { do_read_body(); } else { socket_.close(); timer_.cancel(); } })); } // 修复异步读body的回调逻辑 void do_read_body() { auto self(shared_from_this()); std::cout << "do_read_body\n"; boost::asio::async_read(socket_, boost::asio::buffer(res.body(), res.body_length()), strand_.wrap([this, self](boost::system::error_code ec, std::size_t length) { if (ec) { std::cerr << "read body error:" << ec.value() << " message: " << ec.message() << "\n"; socket_.close(); timer_.cancel(); return; } if (!length) { socket_.close(); timer_.cancel(); return; } if constexpr (log_active) { request_log << "read " << res.body() << "\n"; request_log.flush(); } std::cout << "read " << res.body() << "\n"; request_queue_.send(res.body(), res.body_length(), 0); do_read_header(); // 继续读取下一个消息 })); } // ... 原有成员 boost::asio::steady_timer timer_; // 添加定时器成员 };
方案2:单独线程监听消息队列
如果消息队列的数据非常频繁,定时器轮询可能有延迟,这时候可以创建单独线程阻塞监听消息队列,一旦有数据就post到strand执行写操作。
关键代码修改如下:
class Session : public std::enable_shared_from_this<Session> { public: Session(tcp::socket socket) : socket_(std::move(socket)) , response_queue_(boost::interprocess::open_or_create, "response_queue", 100, 100) , request_queue_(boost::interprocess::open_or_create, "request_queue", 100, 100) , response_listener_(&Session::listen_response_queue, this) // 启动监听线程 { } ~Session() { // 清理资源 socket_.close(); response_queue_.try_send("", 0, 0); // 唤醒阻塞的receive if (response_listener_.joinable()) { response_listener_.join(); } } void start() { post(strand_, [this, self = shared_from_this()] { do_read_header(); }); } private: // 监听消息队列的线程函数 void listen_response_queue() { while (socket_.is_open()) { Message msg; boost::interprocess::message_queue::size_type recvd_size; unsigned int priority = 0; try { response_queue_.receive(msg.body(), msg.max_body_length, recvd_size, priority); // 阻塞等待数据 } catch (const boost::interprocess::interprocess_exception& e) { std::cerr << "message queue receive error: " << e.what() << "\n"; break; } if (!socket_.is_open()) break; // 把写操作post到strand,保证线程安全 post(strand_, [this, self = shared_from_this(), msg = std::move(msg), recvd_size]() mutable { do_write(msg, recvd_size); }); } } void do_write(Message msg, boost::interprocess::message_queue::size_type recvd_size) { try { msg.body_length(recvd_size); msg.encode_header(); std::memcpy(data_, msg.data(), msg.length()); } catch (std::exception& e ) { std::cout << e.what(); return; } std::cout << "write " << msg.body() << "\n"; boost::asio::async_write(socket_, boost::asio::buffer(data_, msg.length()), strand_.wrap([this](boost::system::error_code ec, std::size_t /*length*/) { if (ec) { std::cerr << "write error:" << ec.value() << " message: " << ec.message() << "\n"; socket_.close(); } })); } // ... 原有成员(do_read_header和do_read_body的修复同方案1) std::thread response_listener_; // 添加监听线程成员 };
额外注意事项
- strand的正确使用:所有操作socket的代码都要通过strand执行,避免线程安全问题,上面的代码已经用
strand_.wrap或post(strand_, ...)保证。 - 资源清理:当socket关闭时,要及时停止定时器或join监听线程,避免资源泄漏。
- 异常处理:消息队列的
receive/try_receive可能抛出interprocess_exception,要做好捕获处理。
内容的提问来源于stack exchange,提问作者yaodav
相关产品推荐
相关产品推荐

