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

使用make_parallel_group时WebSocket写入触发I/O中止错误的排查

问题原因分析

  1. WebSocket流操作未通过strand序列化
    WebSocket的stream对象并非线程安全,所有异步操作(async_read、async_write等)必须通过同一个strand调度执行,确保操作的原子性与顺序性。你的代码中co_spawn直接使用IO上下文ioc启动协程,协程可能被调度到线程池的不同线程,违反了WebSocket的线程安全要求,进而触发I/O中止错误。

  2. ParallelGroup取消操作导致流进入错误状态
    使用make_parallel_group::wait_for_one模式时,任一操作完成后,剩余未完成的异步操作会被自动取消(错误码为asio::error::operation_aborted)。如果WebSocket的async_read被取消,流会进入错误状态,后续调用async_write时会直接复用该错误状态,导致写入操作失败。

解决方法

1. 用strand序列化WebSocket操作

修改co_spawn的调用逻辑,使用asio::make_strand(ioc)作为协程的执行上下文,保证所有WebSocket操作都在同一strand上调度:

asio::co_spawn(asio::make_strand(ioc), run_session(std::move(ws), channel), asio::detached);

2. 清除取消操作后的流错误状态

在parallel_group执行完成后,检查WebSocket读取操作的错误码,若为取消错误则清除流的错误状态,恢复流的正常可用状态:

// 在parallel_group结果处理后添加
if (ec2 == asio::error::operation_aborted) {
    ws.clear(ec2);
}

3. 写入前校验连接状态(可选)

发起写入操作前,先检查WebSocket连接是否处于打开状态,避免无效操作:

if (!ws.is_open()) {
    co_return;
}
// 执行写入操作
co_await ws.async_write(buffer, asio::use_awaitable);

修复后的关键代码片段

void run(asio::io_context& ioc, tcp::endpoint endpoint) {
    // ... 其他代码 ...
    // 用strand包装协程执行上下文
    asio::co_spawn(asio::make_strand(ioc), run_session(std::move(ws), channel), asio::detached);
    // ... 其他代码 ...
}

asio::awaitable<void> run_session(websocket::stream<beast::tcp_stream> ws, 
                                  asio::experimental::concurrent_channel<void(asio::error_code, std::string)> channel) {
    // ... 其他代码 ...
    for (;;) {
        auto [idx, ec1, n1, ec2, msg] = co_await make_parallel_group(
            [&](auto token) { return channel.async_receive(token); },
            [&](auto token) { return ws.async_read(buffer, token); }
        ).async_wait(make_parallel_group::wait_for_one, asio::use_awaitable);

        // 清除WebSocket读取操作的取消错误
        if (ec2 == asio::error::operation_aborted) {
            ws.clear(ec2);
        }

        // ... 数据处理逻辑 ...

        // 校验连接状态
        if (!ws.is_open()) {
            co_return;
        }

        // 执行写入操作
        co_await ws.async_write(buffer, asio::use_awaitable);
        // ... 其他代码 ...
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 07:37:13