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

FastAPI WebSocket结合async/await嵌套yield时消息仅结束后发送问题

问题成因

核心问题是同步CPU密集型计算阻塞了AsyncIO事件循环:
FastAPI的WebSocket路由运行在AsyncIO单线程事件循环上,你定义的simulate、simulate_intervals都是同步函数,内部的高负载计算会完全占住事件循环线程。此时即使你调用了await websocket.send_json(),也只是把消息放入了发送队列,事件循环没有空闲时间执行实际的网络IO操作,所有消息都会堆积到所有同步计算执行完毕、事件循环释放控制权后才会批量发送。

解决方案

1. 将CPU密集型计算放到独立的进程池运行

因为是高负载CPU计算,使用进程池可以避开GIL限制,同时不会阻塞AsyncIO事件循环,每次计算完成后事件循环可以立即执行消息发送逻辑。
示例修改代码:

import asyncio
from concurrent.futures import ProcessPoolExecutor
import json

# 全局初始化进程池,可根据CPU核心数调整进程数
executor = ProcessPoolExecutor(max_workers=4)

# 抽离单次计算逻辑,供进程池调用
def run_single_interval(data):
    return interval(data)

def run_single_trial_intervals(data):
    return [interval(data) for _ in range(data.n_intervals)]

@app.websocket("/ws")
async def socket(websocket: WebSocket):
    await websocket.accept()
    loop = asyncio.get_running_loop()
    while True:
        data = await websocket.receive_text()
        nodes = distributions(data)

        nodosJson = json.dumps(nodes, cls=NumpyEncoder)
        await websocket.send_json({"tipo": "nodos", "datos": json.loads(nodosJson)})
        # 发送后主动让出事件循环,确保消息立即发出
        await asyncio.sleep(0)

        # 迭代计算逻辑改为异步调用进程池
        for trialI in range(data.n_trials):
            # 单次trial计算放到进程池执行
            trial = await loop.run_in_executor(executor, run_single_trial_intervals, data)
            for stateI, state in enumerate(trial):
                stateString = json.dumps(state, cls=NumpyEncoder)
                await websocket.send_json(
                    {
                        "tipo": "estado",
                        "datos": json.loads(stateString),
                        "trialI": trialI,
                        "stateI": stateI,
                    }
                )
                # 每次发送后让出事件循环
                await asyncio.sleep(0)

        await websocket.send_json({"tipo": "estado", "msg": "fin"})
        await asyncio.sleep(0)

2. 部署侧配置适配

如果线上部署使用了Nginx等反向代理,需要额外关闭WebSocket代理缓存,避免代理层堆积消息:

location /ws {
    proxy_pass http://你的上游服务地址;
    proxy_http_version 1.1;
    proxy_set_header Upgrade $http_upgrade;
    proxy_set_header Connection "upgrade";
    # 关闭代理缓存
    proxy_buffering off;
    # 调整超时时间适配长耗时任务
    proxy_read_timeout 600s;
    proxy_send_timeout 600s;
}

注意事项

  • 若使用JAX做计算,需确保进程池内的JAX环境初始化正确,避免多进程设备占用冲突,可在进程初始化函数中提前配置JAX运行参数。
  • await asyncio.sleep(0)的作用是主动交回事件循环控制权,强制事件循环处理待发送的网络IO任务,即使设置为0也不会产生明显的性能损耗。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 20:27:03