如何在websockets服务端场景下正确使用asyncio实现无回调的异步等待
解决方案
核心思路
要完全移除回调的需求可以通过asyncio.Future实现,本质是把原来存回调的字典替换为存Future对象,利用Future的可等待特性替代回调触发逻辑。
完整修改后的代码
import asyncio import json from uuid import uuid4 # 字典存储等待响应的Future对象,key为消息id,value为对应的Future实例 PENDING_FUTURES = dict() async def websocket_connection(ws, path): try: async for msg in ws: msg = json.loads(msg) msg_id = msg['id'] msg_data = msg['data'] if msg_id in PENDING_FUTURES: future = PENDING_FUTURES.pop(msg_id) # 给Future设置结果,唤醒对应await的协程 future.set_result(msg_data) else: raise Exception(f"bad id {msg_id}") finally: # 连接断开时清理所有 pending 的Future,避免永久阻塞 for future in PENDING_FUTURES.values(): if not future.done(): future.set_exception(ConnectionError("websocket连接已断开")) PENDING_FUTURES.clear() async def submit(ws, data, timeout=10): request_id = str(uuid4()) # 获取当前运行的事件循环,创建Future对象 loop = asyncio.get_running_loop() future = loop.create_future() PENDING_FUTURES[request_id] = future try: msg = json.dumps({'id': request_id, 'data': data}) await ws.send(msg) # 等待Future完成,加超时避免永久等待 return await asyncio.wait_for(future, timeout=timeout) except Exception as e: # 出现异常时清理字典里的Future if request_id in PENDING_FUTURES: PENDING_FUTURES.pop(request_id) if not future.done(): future.set_exception(e) raise async def websocket_script(ws): # 现在直接await就能拿到返回结果,完全不需要回调 sum = int(await submit(ws, {"eval": "2+2"})) division = float(await submit(ws, {"eval": "3/2"})) print(f"计算结果:2+2={sum}, 3/2={division}") # 连接处理函数,用gather实现两个协程并行运行 async def handle_connection(ws, path): await asyncio.gather( websocket_connection(ws, path), websocket_script(ws) )
逻辑说明
- 并行逻辑确实通过
asyncio.gather实现,它会同时运行传入的多个协程,直到所有协程执行完成或者抛出异常 - Future的作用就是替代回调:当你调用
await submit(...)时,协程会挂起,直到websocket_connection里给对应的Future设置结果,协程才会恢复执行并拿到返回值,全程不需要写任何回调函数 - 代码里加了异常兜底处理:连接断开、请求超时的场景都会主动抛出异常,不会出现协程永久挂死的问题
内容的提问来源于stack exchange,提问作者winwin
相关产品推荐
相关产品推荐

