Python+Starlette:长异步任务运行时如何响应WebSocket消息
解决WebSocket长任务运行时无法响应新消息的问题
你的核心问题在于:收到"Begin"指令后直接await simulate_long_process(websocket),会让当前的WebSocket处理协程被阻塞至长任务结束,即便长任务内部用了await,主协程也无法回到消息循环处理新请求。下面给出两种针对性解决方案,优先满足你关注的纯异步实现需求。
方案1:纯异步实现(基于asyncio Task)
用asyncio.create_task()将长任务包装为后台异步任务,让主协程无需等待长任务完成,立即回到循环处理新的WebSocket消息。
修改后的服务器核心代码
import asyncio from starlette.applications import Starlette from starlette.routing import Route from starlette.endpoints import WebSocketEndpoint async def simulate_long_process(websocket): for i in range(5): await websocket.send_text(f"长任务进度: {i+1}/5") await asyncio.sleep(1) # 主动让出事件循环,允许处理其他任务 await websocket.send_text("长任务完成") class WebSocketTest(WebSocketEndpoint): encoding = "text" async def on_connect(self, websocket): await websocket.accept() async def on_receive(self, websocket, data): if data == "Begin": # 创建后台任务,无需await,长任务将独立并发执行 asyncio.create_task(simulate_long_process(websocket)) await websocket.send_text("长任务已启动") elif data.startswith("Send"): # 即时处理Send请求,无需等待长任务 random_num = data.split(":")[-1] await websocket.send_text(f"收到随机数: {random_num}") app = Starlette(routes=[Route("/", WebSocketTest)])
原理说明
asyncio.create_task()会把长任务协程注册到事件循环,作为独立任务运行,主协程不会被阻塞。- 长任务内部的
await asyncio.sleep()会主动让出事件循环,此时事件循环可以处理WebSocket的新消息接收与回复。
如果需要更安全地管理后台任务(比如连接关闭时自动取消未完成的长任务),可以用asyncio.TaskGroup():
class WebSocketTest(WebSocketEndpoint): encoding = "text" _task_group = None async def on_connect(self, websocket): await websocket.accept() self._task_group = asyncio.TaskGroup() self._task_group.__enter__() async def on_receive(self, websocket, data): if data == "Begin": self._task_group.create_task(simulate_long_process(websocket)) await websocket.send_text("长任务已启动") elif data.startswith("Send"): random_num = data.split(":")[-1] await websocket.send_text(f"收到随机数: {random_num}") async def on_disconnect(self, websocket, close_code): if self._task_group: await self._task_group.__aexit__(None, None, None)
方案2:线程方案(适用于CPU密集型长任务)
如果你的长任务是CPU密集型(而非IO/等待型),异步协程无法让出事件循环,此时可以用asyncio.to_thread()将任务放到独立线程运行:
# 假设simulate_long_process是CPU密集型同步函数 def simulate_long_process(websocket): import time loop = asyncio.get_event_loop() for i in range(5): # 线程中需通过run_coroutine_threadsafe提交异步操作到主事件循环 asyncio.run_coroutine_threadsafe(websocket.send_text(f"长任务进度: {i+1}/5"), loop) time.sleep(1) asyncio.run_coroutine_threadsafe(websocket.send_text("长任务完成"), loop) class WebSocketTest(WebSocketEndpoint): encoding = "text" async def on_receive(self, websocket, data): if data == "Begin": await asyncio.to_thread(simulate_long_process, websocket) await websocket.send_text("长任务已启动") elif data.startswith("Send"): random_num = data.split(":")[-1] await websocket.send_text(f"收到随机数: {random_num}")
注意事项
- 线程内不能直接执行异步操作,必须用
asyncio.run_coroutine_threadsafe()将WebSocket消息发送操作提交到主事件循环。 - 此方案仅适合CPU密集型任务,IO密集型任务优先选择纯异步方案。
内容的提问来源于stack exchange,提问作者rob_7cc
相关产品推荐
相关产品推荐

