Sanic框架WebSocket接口中await ws.send消息未即时发送问题咨询
问题根因分析
- Python异步 runtime 基于单线程事件循环模型,所有协程调度、IO操作的执行都依赖事件循环的空闲时间片,你直接在协程内部调用未做异步改造的长耗时同步函数
process_long_task时,会完全占用当前线程,导致整个事件循环被阻塞,无法执行任何其他调度任务。 - Sanic的
await ws.send调用只会将消息写入到Socket发送缓冲区,不会等待内核实际完成网络发送操作就会切换回当前协程。你调用完发送方法后立刻进入阻塞的同步任务,事件循环没有机会调度底层的网络IO把缓冲区的第一条消息推送给客户端,直到同步任务执行完毕、事件循环恢复调度后,才会把两条积压的消息一起发出。 - 你添加
await asyncio.sleep(0.1)能临时生效的本质是:await动作会主动让出当前协程的执行权,给事件循环留出时间处理完积压的网络发送任务,把第一条消息发出去之后再回到当前协程执行同步任务。但这个方案只是不稳定的 workaround,如果事件队列中存在其他待处理任务,0.1秒的等待不一定能保证消息发送完成。 - 额外说明:直接在协程中执行长同步任务不仅会导致当前WebSocket消息延迟,还会阻塞整个Sanic服务的所有其他请求(包括HTTP请求、其他WebSocket连接的消息处理),严重影响服务并发能力。
正确解决方案
根据你的同步任务类型选择对应方案即可,无需额外添加sleep逻辑就能保证消息实时发送:
方案1:IO密集型同步任务用线程池执行
如果你的process_long_task是磁盘IO、网络请求等IO密集型任务,用asyncio.to_thread扔到独立线程执行即可,不会阻塞事件循环:
import json import asyncio from sanic import Sanic from sanic.server.websockets.connection import WebSocketConnection app = Sanic("TestApp") # 无法改造的长耗时同步任务 def process_long_task(data): # 业务逻辑实现 pass @app.websocket("/") async def process_task(request, ws: WebSocketConnection): raw_data = await ws.recv() data = json.loads(raw_data) await ws.send("TASK STARTED") # 将同步任务提交到独立线程执行,当前协程等待任务完成后继续 await asyncio.to_thread(process_long_task, data) await ws.send("TASK ENDED")
方案2:CPU密集型同步任务用进程池执行
如果你的process_long_task是计算密集型任务,为了避免GIL影响性能,用进程池执行更合适:
import json import asyncio from concurrent.futures import ProcessPoolExecutor from sanic import Sanic from sanic.server.websockets.connection import WebSocketConnection app = Sanic("TestApp") # 全局进程池实例,可根据服务器配置调整进程数量 executor = ProcessPoolExecutor(max_workers=4) def process_long_task(data): # 计算密集型业务逻辑实现 pass @app.websocket("/") async def process_task(request, ws: WebSocketConnection): raw_data = await ws.recv() data = json.loads(raw_data) await ws.send("TASK STARTED") # 将同步任务提交到进程池执行 await asyncio.get_event_loop().run_in_executor(executor, process_long_task, data) await ws.send("TASK ENDED")
内容的提问来源于stack exchange,提问作者Инт Альгамбра
相关产品推荐
相关产品推荐

