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

WebSocket子协程中co_await调用出现阻塞问题排查

问题:WebSocket协程复用对象时发送逻辑阻塞

我尝试在两个协程中复用同一个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
问题原因与修复方案

核心问题分析

  1. 错误创建独立io_context
    代码中在do_session里新建了ioc2,但WebSocket对象ws绑定的是main函数中ioc的执行器。将协程提交到ioc2后,ioc2.run()会阻塞,而WebSocket的IO事件只会在原ioc的线程中处理,导致async_write/async_read的完成回调无法触发,协程永久挂起。

  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 19:49:51