boost::beast::websocket::stream异步读写实现方式咨询
boost::beast::websocket::stream异步读写实现方式咨询
嗨,刚好对Boost.Beast WebSocket的异步读写有不少实践经验,来给你详细讲讲~
首先回答你的核心问题:完全可以继续用std::queue<std::string>来存储待发送的消息,不需要换成flat_buffer,适配起来和你之前写socket/SSL流的队列模式逻辑几乎一致,只需要针对WebSocket的接口做少量调整。下面给你拆解异步写和异步读的实现细节,再附上完整的示例代码。
一、异步写的实现思路
WebSocket的async_write接口支持直接接收std::string的buffer,所以你的std::queue<std::string>队列可以直接复用,核心逻辑还是保证同一时刻只有一个异步写操作在执行(避免并发写冲突):
- 调用
write方法时,把消息加入队列; - 如果当前没有正在执行的写操作,就启动第一个异步写;
- 每一个异步写完成后,弹出队列头部的消息,若队列非空则继续启动下一个异步写。
需要注意的小细节:WebSocket发送的是帧,你可以在构造时或写操作前指定发送的是文本帧还是二进制帧(比如m_ws.text(true)设置为文本帧)。
二、异步读的实现思路
关于读操作,推荐用boost::beast::flat_buffer作为类成员,原因如下:
flat_buffer是内存连续的缓冲区,读取效率比分散缓冲区更高;- 不需要每次读完都手动
clear:Beast的异步读接口会自动把新数据追加到缓冲区的可用区域,你处理完数据后,只需要调用consume(bytes_transferred)释放已经处理的内存即可,缓冲区会自动复用剩余空间,避免频繁内存分配。如果强制clear反而会浪费之前分配的内存,降低性能。
三、完整的异步读写类示例
#include <boost/beast.hpp> #include <boost/asio.hpp> #include <queue> #include <string> #include <memory> #include <mutex> // 如果需要多线程安全的话添加 namespace beast = boost::beast; namespace asio = boost::asio; using tcp = asio::ip::tcp; using websocket_stream = beast::websocket::stream<tcp::socket>; template <typename WsStream> class WsSession : public std::enable_shared_from_this<WsSession<WsStream>> { public: explicit WsSession(WsStream ws_stream) : m_ws(std::move(ws_stream)) { // 默认设置为文本帧,若需要二进制帧可改为m_ws.binary(true) m_ws.text(true); } void start() { // 启动第一个异步读操作 do_read(); } // 线程安全的写方法(单线程IO上下文可省略锁) void write(std::string data) { std::lock_guard<std::mutex> lock(m_write_mutex); // 多线程环境下加锁 bool write_in_progress = !m_write_queue.empty(); m_write_queue.push(std::move(data)); if (!write_in_progress) { do_write(); } } void close() { auto self = this->shared_from_this(); // 异步关闭WebSocket连接,使用正常关闭码 m_ws.async_close(beast::websocket::close_code::normal, [self](beast::error_code ec) { if (ec && ec != beast::websocket::error::closed) { // 处理关闭过程中的错误 } }); } private: WsStream m_ws; beast::flat_buffer m_read_buffer; // 复用的读缓冲区 std::queue<std::string> m_write_queue; std::mutex m_write_mutex; // 多线程环境下保护队列的锁 void do_read() { auto self = this->shared_from_this(); // 异步读取WebSocket帧 m_ws.async_read(m_read_buffer, [this, self](beast::error_code ec, std::size_t bytes_transferred) { if (!ec) { // 把缓冲区的数据转换为字符串处理 std::string received_data = beast::buffers_to_string(m_read_buffer.data()); // 释放已经处理的缓冲区空间,方便后续复用 m_read_buffer.consume(bytes_transferred); // 继续启动下一个异步读 do_read(); } else if (ec != beast::websocket::error::closed) { // 处理非连接关闭的读错误 } }); } void do_write() { auto self = this->shared_from_this(); std::lock_guard<std::mutex> lock(m_write_mutex); // 多线程环境下加锁 if (m_write_queue.empty()) return; const auto& current_msg = m_write_queue.front(); // 异步发送WebSocket帧 m_ws.async_write(asio::buffer(current_msg), [this, self](beast::error_code ec, std::size_t /*bytes_sent*/) { std::lock_guard<std::mutex> lock(m_write_mutex); if (!ec) { m_write_queue.pop(); if (!m_write_queue.empty()) { // 队列还有消息,继续启动下一个异步写 do_write(); } } else if (ec != beast::websocket::error::closed) { // 处理写错误,可根据情况清空队列或做其他处理 m_write_queue = {}; } }); } };
关键细节补充
- 线程安全:如果你的
write方法会被多个线程调用,一定要给m_write_queue加锁(比如示例中的std::mutex),因为std::queue本身不是线程安全的。如果所有操作都在同一个io_context的线程中执行,就不需要加锁。 - 帧类型切换:如果需要混合发送文本和二进制帧,可以在
write方法中添加参数指定帧类型,在do_write前临时设置m_ws.text(true)或m_ws.binary(true)。 - 错误处理:示例中简化了错误处理逻辑,实际项目中建议根据具体错误码做对应的处理(比如连接断开时清理资源、记录日志等)。
内容来源于stack exchange
相关产品推荐
相关产品推荐

