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

WebSocket进度消息在流程结束时批量返回而非逐条推送问题排查

问题排查结论

你遇到的消息批量返回问题核心是异步事件循环被同步阻塞代码卡死,附带两个代码缺陷:

  • 核心故障:create_mds、create_sample_data、correlation_matrix、accomplish_mds这几个都是同步执行的CPU密集型函数,你在async协程里直接调用这些函数时,当前线程会被完全占用来跑计算,asyncio事件循环根本没机会调度执行websocket的网络发送逻辑,所有进度消息都会积压在内存队列里,等全部计算跑完、事件循环拿回控制权,才会把攒的所有消息一次性发出去。第一条消息能正常出现,是因为发送完第一条之后才进入第一个阻塞函数,发送操作刚好赶在阻塞前被调度了。
  • 缺陷1:ConnectionManager.send_response的异常分支里调用disconnect没加await,disconnect本身是异步协程,不加await不会真正执行连接关闭逻辑,还会抛出未等待协程的运行时警告。
  • 缺陷2:当前所有计算任务直接跑在websocket路由的协程里,只要有一个用户的计算任务在跑,同进程下其他所有用户的websocket连接都会被连带阻塞,无法收发消息。
修复方法
  1. 将同步阻塞的计算逻辑挪到事件循环外执行,避免卡住消息发送
    用Python 3.9+自带的asyncio.to_thread把同步函数包装为非阻塞调用,每次计算时主动让出事件循环,保证进度消息能即时推送。修改后的automate_algorithm代码如下:
    import asyncio
    from functools import partial
    
    async def automate_algorithm(idt,language='en'):
        from controllers.dispatcher_controller import send_message_to_socket
        
        await send_message_to_socket(client_id=idt, what="process starting")
        
        # 同步函数扔到独立线程执行,不阻塞事件循环
        mds2 = await asyncio.to_thread(partial(create_mds, idt, mapper=False))
        await send_message_to_socket(client_id=idt, what="main_data_structure 2 created...")
        
        sample_data = await asyncio.to_thread(partial(create_sample_data, idt, mapper=False))
        await send_message_to_socket(client_id=idt, what="sample data created...")
        
        corr = await asyncio.to_thread(partial(correlation_matrix, idt, mapper=False))
        await send_message_to_socket(client_id=idt, what="correlation created...")
        
        mds3 = await asyncio.to_thread(partial(accomplish_mds, idt, mapper=False))
        await send_message_to_socket(client_id=idt, what="main_data_structure 3 created...", status="completed")
    

    注意:如果你的计算逻辑是纯CPU密集型,线程池会受Python GIL限制仍有轻微卡顿,可以替换为asyncio.to_thread对应的进程池实现,彻底避免阻塞。不要试图在同步计算函数里加asyncio.sleep(0)让出控制权,这种方式对纯同步的CPU计算代码无效。

  2. 修复异步调用缺失的问题
    找到ConnectionManager.send_response里的异常捕获块,给disconnect调用加上await:
    except Exception as e:
        print(str(e))
        await self.disconnect(id)
    
  3. 可选优化:将计算任务从websocket协程中解耦,用后台任务队列独立执行,进一步降低单任务阻塞全站的风险。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 06:30:50