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

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>队列可以直接复用,核心逻辑还是保证同一时刻只有一个异步写操作在执行(避免并发写冲突):

  1. 调用write方法时,把消息加入队列;
  2. 如果当前没有正在执行的写操作,就启动第一个异步写;
  3. 每一个异步写完成后,弹出队列头部的消息,若队列非空则继续启动下一个异步写。

需要注意的小细节: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 = {};
                }
            });
    }
};

关键细节补充

  1. 线程安全:如果你的write方法会被多个线程调用,一定要给m_write_queue加锁(比如示例中的std::mutex),因为std::queue本身不是线程安全的。如果所有操作都在同一个io_context的线程中执行,就不需要加锁。
  2. 帧类型切换:如果需要混合发送文本和二进制帧,可以在write方法中添加参数指定帧类型,在do_write前临时设置m_ws.text(true)或m_ws.binary(true)。
  3. 错误处理:示例中简化了错误处理逻辑,实际项目中建议根据具体错误码做对应的处理(比如连接断开时清理资源、记录日志等)。

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.08 08:58:06