WebSocket子协程中co_await调用出现阻塞问题排查
我尝试在两个协程中复用同一个WebSocket对象,分别实现发送和接收逻辑。但在已有协程中调用子协程时,代码在co_await async_write处发生阻塞。父协程中直接用while(1)循环调用async_write能正常运行,但子协程仅能成功执行一次,后续就会阻塞。想问调用子协程的方式是不是有问题?
客户端代码
// // Copyright (c) 2022 Klemens D. Morgenstern (klemens dot morgenstern at gmx dot net) // // Distributed under the Boost Software License, Version 1.0. (See accompanying // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt) // // Official repository: https://github.com/boostorg/beast // //------------------------------------------------------------------------------ // // Example: WebSocket client, coroutine // //------------------------------------------------------------------------------ #include <boost/beast/core.hpp> #include <boost/beast/websocket.hpp> #include <cstdlib> #include <functional> #include <iostream> #include <string> #include <boost/asio/awaitable.hpp> #include <boost/asio/co_spawn.hpp> #include <boost/asio/detached.hpp> #include <boost/asio/use_awaitable.hpp> #include <thread> #include <chrono> #if defined(BOOST_ASIO_HAS_CO_AWAIT) namespace beast = boost::beast; // from <boost/beast.hpp> namespace http = beast::http; // from <boost/beast/http.hpp> namespace websocket = beast::websocket; // from <boost/beast/websocket.hpp> namespace net = boost::asio; // from <boost/asio.hpp> using tcp = boost::asio::ip::tcp; // from <boost/asio/ip/tcp.hpp> //------------------------------------------------------------------------------ // Sends a WebSocket message and prints the response net::awaitable<void> do_write_loop(auto& ws, std::string text) { std::cout << "before while"; //std::cout << "do_write_loop"<<i<<std::endl; co_await ws.async_write(net::buffer(std::string(text))); } net::awaitable<void> do_read_loop(auto& ws) { beast::flat_buffer buffer; while (true) { std::cout << "do_read_loop"; co_await ws.async_read(buffer); std::cout << beast::make_printable(buffer.data()) << std::endl; } } net::awaitable<void> do_session( std::string host, std::string port, std::string text) { auto resolver = net::use_awaitable.as_default_on( tcp::resolver(co_await net::this_coro::executor)); auto ws = net::use_awaitable.as_default_on( websocket::stream<beast::tcp_stream>(co_await net::this_coro::executor)); auto const results = co_await resolver.async_resolve(host, port); beast::get_lowest_layer(ws).expires_after(std::chrono::seconds(30)); auto ep = co_await beast::get_lowest_layer(ws).async_connect(results); host += ':' + std::to_string(ep.port()); beast::get_lowest_layer(ws).expires_never(); ws.set_option(websocket::stream_base::timeout::suggested(beast::role_type::client)); ws.set_option(websocket::stream_base::decorator( [](websocket::request_type& req) { req.set(http::field::user_agent, std::string(BOOST_BEAST_VERSION_STRING) + " websocket-client-coro"); })); co_await ws.async_handshake(host, "/"); std::cout << "after handshake"; //while (1) { // co_await ws.async_write(net::buffer(std::string(text))); //} net::io_context ioc2; auto executor = co_await boost::asio::this_coro::executor; // Start the write and read loops concurrently net::co_spawn(ioc2, do_write_loop(ws, text), [](std::exception_ptr e) { if (e) try { std::cout << "enter e"; std::rethrow_exception(e); } catch (std::exception& e) { std::cerr << "Error: " << e.what() << "\n"; } }); net::co_spawn(ioc2, do_read_loop(ws), [](std::exception_ptr e) { if (e) try { std::rethrow_exception(e); } catch (std::exception& e) { std::cerr << "Error: " << e.what() << "\n"; } }); ioc2.run(); //ioc3.run(); // Wait for both tasks to complete (this will never happen) //co_await net::when_all(std::move(write_loop_task), std::move(read_loop_task)); // Close the WebSocket connection (this will never be reached) co_await ws.async_close(websocket::close_code::normal); } //------------------------------------------------------------------------------ int main(int argc, char **argv) { // Check command line arguments. /*if (argc != 4) { std::cerr << "Usage: websocket-client-awaitable <host> <port> <text>\n" << "Example:\n" << " websocket-client-awaitable echo.websocket.org 80 \"Hello, world!\"\n"; return EXIT_FAILURE; }*/ //auto const host = argv[1]; //auto const port = argv[2]; //auto const text = argv[3]; auto const host = "127.0.0.1"; auto const port = "12345"; auto const text = "argv[3]"; // The io_context is required for all I/O net::io_context ioc; // Launch the asynchronous operation net::co_spawn(ioc, do_session(host, port, text), [](std::exception_ptr e) { if (e) try { std::rethrow_exception(e); } catch (std::exception &e) { std::cerr << "Error: " << e.what() << "\n"; } }); // Run the I/O service. The call will return when // the socket is closed. ioc.run(); return EXIT_SUCCESS; } #else int main(int, char *[]) { std::printf("awaitables require C++20\n"); return 1; } #endif
核心问题分析
错误创建独立
io_context
代码中在do_session里新建了ioc2,但WebSocket对象ws绑定的是main函数中ioc的执行器。将协程提交到ioc2后,ioc2.run()会阻塞,而WebSocket的IO事件只会在原ioc的线程中处理,导致async_write/async_read的完成回调无法触发,协程永久挂起。do_write_loop缺少循环逻辑
当前do_write_loop仅执行一次async_write就结束协程,这就是发送逻辑只执行一次的直接原因,和阻塞无关,但也是功能缺失的问题。
修复步骤
1. 复用原执行器,删除独立io_context
直接使用co_await net::this_coro::executor获取的原执行器启动协程,不需要新建ioc2。所有IO操作都在同一个io_context中处理,确保回调能正常触发。
2. 给do_write_loop添加循环逻辑
要实现持续发送,必须在do_write_loop中添加while(true)循环,同时建议添加时间间隔,避免无限制发送导致网络压力过大。
3. 正确管理协程生命周期
使用net::detached让子协程后台运行,同时确保do_session协程不会提前结束(否则ws对象会被销毁,子协程会崩溃),可以通过无限等待的定时器实现。
修复后的关键代码
// 修改后的do_write_loop,添加循环和间隔 net::awaitable<void> do_write_loop(auto& ws, std::string text) { while (true) { std::cout << "发送消息中\n"; co_await ws.async_write(net::buffer(text)); // 每秒发送一次,可根据需求调整 co_await net::steady_timer(co_await net::this_coro::executor, std::chrono::seconds(1)).async_wait(net::use_awaitable); } } // 修改后的do_session协程启动部分 net::awaitable<void> do_session( std::string host, std::string port, std::string text) { // ... 前面的握手逻辑保持不变 ... std::cout << "握手完成\n"; // 在原执行器上启动子协程,用detached让它们后台运行 net::co_spawn(co_await net::this_coro::executor, do_write_loop(ws, text), net::detached); net::co_spawn(co_await net::this_coro::executor, do_read_loop(ws), net::detached); // 无限等待,保持do_session协程存活,避免ws对象被销毁 co_await net::steady_timer(co_await net::this_coro::executor, std::chrono::steady_clock::time_point::max()).async_wait(net::use_awaitable); // 实际场景中可在此处处理WebSocket关闭逻辑 // co_await ws.async_close(websocket::close_code::normal); }
额外注意事项
- 线程安全:Boost Beast的
websocket::stream并非线程安全,但在同一个执行器(即同一个线程,或通过strand序列化)上的协程调用是安全的,所有IO操作都在原ioc的线程中处理,不会出现并发访问问题。 - 错误处理:建议在子协程内部添加异常捕获,比如捕获
async_write/async_read抛出的异常,避免单个协程崩溃导致整个程序退出。 - 资源释放:实际应用中需要添加退出逻辑(比如接收退出信号),不要一直无限等待,确保资源能正常释放。
内容的提问来源于stack exchange,提问作者user24267314

