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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.23 18:36:04