如何在Boost WebSocket异步服务器中实现async_read与async_write独立
Boost WebSocket 异步读写独立实现问题
你的当前实现存在的问题
- 同步
ws_.write()会阻塞当前线程,破坏Boost.Asio的异步事件循环机制,导致其他异步操作(比如新客户端连接、其他会话的读写)被延迟处理。 - 在
do_read()中触发do_write()的逻辑会让读写操作串行化,无法真正实现独立运行;如果写入操作耗时较长,还会直接阻塞读操作的响应。 - 仅靠
write_flag标记无法处理并发写入场景,多次触发写操作可能导致消息丢失或重复发送。
正确的异步读写独立实现方案
要实现异步读、写操作完全独立运行,核心要做到两点:
- 保持异步读操作持续运行:一个读操作完成后立即发起下一个,确保随时能接收客户端消息。
- 维护待发送消息队列:主动发消息时将消息加入队列;若当前无正在执行的异步写操作,就从队列取消息发起异步写;写操作完成后,检查队列是否还有消息,有则继续发起下一个写。
示例代码实现:
#include <boost/beast.hpp> #include <boost/asio.hpp> #include <queue> #include <memory> #include <string> namespace beast = boost::beast; namespace websocket = beast::websocket; namespace net = boost::asio; using tcp = boost::asio::ip::tcp; class session : public std::enable_shared_from_this<session> { public: explicit session(tcp::socket socket) : ws_(std::move(socket)) {} void run() { ws_.set_option(websocket::stream_base::timeout::suggested(beast::role_type::server)); ws_.set_option(websocket::stream_base::decorator( [](websocket::response_type& res) { res.set(beast::http::field::server, std::string(BOOST_BEAST_VERSION_STRING) + " websocket-server-async"); })); ws_.async_accept(beast::bind_front_handler(&session::on_accept, shared_from_this())); } // 外部主动发送消息的接口 void send(std::string message) { // 确保操作在WebSocket对应的io_context线程中执行,避免线程安全问题 net::post(ws_.get_executor(), [self = shared_from_this(), msg = std::move(message)]() { bool write_in_progress = !self->write_queue_.empty(); self->write_queue_.push(std::move(msg)); if (!write_in_progress) { self->do_write(); } }); } private: websocket::stream<tcp::socket> ws_; beast::flat_buffer read_buffer_; std::queue<std::string> write_queue_; void on_accept(beast::error_code ec) { if (ec) return fail(ec, "accept"); do_read(); } void do_read() { ws_.async_read(read_buffer_, beast::bind_front_handler(&session::on_read, shared_from_this())); } void on_read(beast::error_code ec, std::size_t bytes_transferred) { boost::ignore_unused(bytes_transferred); if (ec == websocket::error::closed) return; if (ec) return fail(ec, "read"); // 处理读取到的客户端消息(按需实现) std::string received(beast::buffers_to_string(read_buffer_.data())); read_buffer_.consume(read_buffer_.size()); // 立即发起下一个读操作,保持读通道持续可用 do_read(); } void do_write() { ws_.async_write(net::buffer(write_queue_.front()), beast::bind_front_handler(&session::on_write, shared_from_this())); } void on_write(beast::error_code ec, std::size_t bytes_transferred) { boost::ignore_unused(bytes_transferred); if (ec) return fail(ec, "write"); write_queue_.pop(); if (!write_queue_.empty()) { // 队列还有未发消息,继续发起下一个写操作 do_write(); } } void fail(beast::error_code ec, char const* what) { // 错误处理逻辑(比如打印日志、关闭会话等) std::cerr << what << ": " << ec.message() << std::endl; } };
关键细节说明
- 消息队列
write_queue_:缓存待发送消息,避免同时发起多个异步写操作(Boost Beast WebSocket不允许存在未完成的并行异步写操作)。 net::post():确保发送消息的操作在WebSocket所属的io_context线程中执行,因为WebSocket对象并非线程安全,所有操作必须在同一执行上下文(线程)中进行。- 持续异步读:
on_read完成后立即调用do_read(),保证客户端消息能被及时接收,不受写操作影响。 - 独立写触发逻辑:外部调用
send()即可主动发消息,无需依赖读操作触发,写操作会根据队列状态自动连续执行。
总结
你的原方案存在阻塞线程、无法真正实现读写独立的问题,采用消息队列+异步写串行化的方式,既符合Boost.Asio的异步编程模型,又能保证读写操作独立运行,同时避免线程安全风险和阻塞问题。
内容的提问来源于stack exchange,提问作者Srinidhi
相关产品推荐
相关产品推荐

