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

如何在Boost WebSocket异步服务器中实现async_read与async_write独立

Boost WebSocket 异步读写独立实现问题

你的当前实现存在的问题

  • 同步ws_.write()会阻塞当前线程,破坏Boost.Asio的异步事件循环机制,导致其他异步操作(比如新客户端连接、其他会话的读写)被延迟处理。
  • 在do_read()中触发do_write()的逻辑会让读写操作串行化,无法真正实现独立运行;如果写入操作耗时较长,还会直接阻塞读操作的响应。
  • 仅靠write_flag标记无法处理并发写入场景,多次触发写操作可能导致消息丢失或重复发送。

正确的异步读写独立实现方案

要实现异步读、写操作完全独立运行,核心要做到两点:

  1. 保持异步读操作持续运行:一个读操作完成后立即发起下一个,确保随时能接收客户端消息。
  2. 维护待发送消息队列:主动发消息时将消息加入队列;若当前无正在执行的异步写操作,就从队列取消息发起异步写;写操作完成后,检查队列是否还有消息,有则继续发起下一个写。

示例代码实现:

#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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 09:30:53