使用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
相关产品推荐
相关产品推荐

