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

使用Boost Asio和Beast协程调用async_write发送频繁大请求时崩溃

问题分析与解决方案

崩溃原因

崩溃的核心是Beast Websocket Stream不允许同时发起多个同类型异步操作(比如并发调用async_write)。虽然你使用了同一个strand,但strand的串行执行特性无法阻止这种违规调用:

  • 频繁触发handleRequest时,每个sendRequest协程会被调度到strand上执行。
  • 第一个协程执行到co_await ws_->async_write时,会释放strand的执行权(awaitable会让出控制权直到操作完成),strand会立刻调度下一个sendRequest协程。
  • 第二个协程直接再次调用ws_->async_write,此时前一个async_write还未完成,触发Beast内部soft_mutex的断言检查(BOOST_ASSERT(id_ != T::id)),导致程序崩溃。

同步ws_->write能正常运行,是因为同步调用会阻塞直到写入完成,天然保证了同一时间只有一个写入操作在进行。


解决方案

方案1:使用可等待互斥锁(Awaitable Mutex)

利用Asio的experimental::awaitable_mutex保证同一时间只有一个协程执行async_write,避免并发调用:

class Adapter {
public:
    Adapter(asio::strand<boost::asio::io_context::executor_type>& strand)
        : strand_(strand)
        , ssl_context_(asio::ssl::context::tlsv12_client)
        , ws_(new websocket::stream<beast::ssl_stream<beast::tcp_stream>>(strand, ssl_context_))
        , write_mutex_() {}

    void handleRequest(std::string& data) {
        // 处理得到new_data
        co_spawn(strand_, sendRequest(new_data), boost::asio::detached);
    }

    asio::awaitable<void> sendRequest(const std::string& data) {
        // 数据转换得到new_data
        co_await write_mutex_.async_lock(asio::use_awaitable);
        try {
            co_await ws_->async_write(strand_, asio::buffer(new_data), asio::use_awaitable);
        } finally {
            write_mutex_.unlock();
        }
        co_return;
    }

protected:
    asio::strand<boost::asio::io_context::executor_type>& strand_;
    asio::ssl::context  ssl_context_;
    std::unique_ptr<websocket::stream<beast::ssl_stream<beast::tcp_stream>>> ws_;
    asio::experimental::awaitable_mutex write_mutex_; // 添加可等待互斥锁
};

方案2:使用发送队列(适合高频率请求场景)

维护一个发送队列,用单独的协程负责从队列取数据并执行写入,彻底避免并发写入,同时平滑高频率请求流量:

class Adapter {
public:
    Adapter(asio::strand<boost::asio::io_context::executor_type>& strand)
        : strand_(strand)
        , ssl_context_(asio::ssl::context::tlsv12_client)
        , ws_(new websocket::stream<beast::ssl_stream<beast::tcp_stream>>(strand, ssl_context_))
        , is_sending_(false) {}

    void handleRequest(std::string& data) {
        // 处理得到new_data
        {
            std::lock_guard<std::mutex> lock(queue_mutex_);
            send_queue_.push(std::move(new_data));
        }
        // 仅当无发送任务时启动发送协程
        if (!is_sending_.exchange(true)) {
            co_spawn(strand_, processSendQueue(), boost::asio::detached);
        }
    }

    asio::awaitable<void> processSendQueue() {
        try {
            while (true) {
                std::string data;
                {
                    std::lock_guard<std::mutex> lock(queue_mutex_);
                    if (send_queue_.empty()) break;
                    data = std::move(send_queue_.front());
                    send_queue_.pop();
                }
                // 数据转换得到new_data
                co_await ws_->async_write(strand_, asio::buffer(new_data), asio::use_awaitable);
            }
        } finally {
            is_sending_ = false;
        }
        co_return;
    }

protected:
    asio::strand<boost::asio::io_context::executor_type>& strand_;
    asio::ssl::context  ssl_context_;
    std::unique_ptr<websocket::stream<beast::ssl_stream<beast::tcp_stream>>> ws_;
    std::queue<std::string> send_queue_;
    std::mutex queue_mutex_;
    std::atomic<bool> is_sending_;
};

方案3:改用同步写入(简单但可能影响性能)

如果业务场景对性能要求不高,直接使用ws_->write即可,这也是你测试中能正常运行的方式,但高频率大流量场景下可能会阻塞strand,影响其他操作执行。


内容的提问来源于stack exchange,提问作者Walk Anywhere

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 05:25:55