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

async.Event已设置但仍卡在await等待状态的技术问题排查

问题:WebSocket通道已打开但asyncio.Event.wait()始终阻塞

我使用Python的websockets库创建WebSocket连接,后台启动无限任务监听服务器消息并转发至处理器。为订阅实时报价,先发送通道开启请求,用asyncio.Event跟踪通道状态。收到服务器CHANNEL_OPENED消息时,消息处理器中已设置该event,但主程序始终卡在await event.wait()处无法推进。日志显示WebSocket连接成功且通道已打开,断点调试可见event已被设置,但程序一直陷入心跳包循环,无法执行后续代码。

环境信息:

  • Python 3.12.4
  • websockets 12.0

主程序核心代码

import asyncio
from api import ApiClient

async def main():
    api = ApiClient(
        uid=_uid,
        remember_token=_remember_token
    )

    asyncio.create_task(api.message_handler())

    api.request_channel(1)  # 发送通道开启请求,在api.open_channels中创建对应event

    event = api.open_channels.get(channel)
    await event.wait()  # 卡在此处
    log.info("Subscribing options quotes for tracked symbols.")  # 永远无法执行到这行
    for symbol in symbols:  # symbols是字符串列表
        response = api.get(f"/option-chains/{symbol}")
        options = [OptionInstrument(**option) for option in response.json()["data"]["items"]]

        for option in options:
            subscription = asyncio.create_task(api.subscribe(option, channel))
            await subscription

if __name__ == "__main__":
    asyncio.run(main())

通道请求逻辑

def request_channel(self, channel: int):
    if self.open_channels.get(channel):
        log.info(f"Channel {channel} was already requested. Skipping request.")
        return

    self.open_channels[channel] = asyncio.Event()
    self.streamer.send(
        dumps(
            {
                "type": "CHANNEL_REQUEST",
                "channel": channel,
                "service": "FEED",
                "parameters": {"contract": "AUTO"},
            }
        )
    )
    log.debug(f"Requested opening of channel {channel}")

ApiClient消息处理循环

async def message_handler(self):
    while True:
        await self._process_message(self.streamer.recv())

async def _process_message(self, message: str):
    message = loads(message)
    log.debug(message)
    match message["type"]:
        case "AUTH_STATE" | "SETUP":
            if message.get("state") == "AUTHORIZED":
                log.debug("Websocket streamer setup and authorized by the server.")

        case "CHANNEL_OPENED":
            self.open_channels[message["channel"]].set()
            log.debug(f"Channel {message['channel']} open.")

        case "KEEPALIVE":
            self.streamer.send(dumps({"type": "KEEPALIVE", "channel": message["channel"]}))
            log.debug("Extended keep alive")

        case "ERROR":
            log.error(f"Streamer had an error. Restarting. Error message:\n\t{message}")
            self.streamer.close()
            self.open_channels = dict()
            self._streamer = None

        case _:
            log.warning(f"Unexpected message type received:\n\t{message}")

WebSocket流对象配置

@property
def streamer(self) -> ClientConnection:
    if not self._streamer:
        self._streamer = self._setup_streamer_websocket()
    return self._streamer

def _setup_streamer_websocket(self) -> ClientConnection:
    websocket = connect(self.quote_streamer_url)

    setup_message_payload = {
        "type": "SETUP",
        "channel": 0,
        "keepaliveTimeout": 60,
        "acceptKeepaliveTimeout": 60,
        "version": "0.1",
    }

    websocket.send(message=dumps(setup_message_payload))

    auth_message_payload = {
        "type": "AUTH",
        "channel": 0,
        "token": self.quote_streamer_token,
    }

    websocket.send(dumps(auth_message_payload))

    return websocket

问题根源分析

  1. WebSocket异步API被同步调用:websockets库的connect()、send()、recv()都是异步方法,当前代码直接同步调用这些方法,得到的是协程对象而非实际连接/结果,导致消息收发逻辑完全失效。
  2. 消息处理循环的异步调用错误:message_handler中直接传递self.streamer.recv()给_process_message,但recv()是异步方法,必须用await获取结果后再传入。
  3. 事件循环调度异常:错误的同步操作打乱了事件循环的任务调度,即使event.set()被执行,事件循环也无法及时唤醒等待event.wait()的任务。

修复方案

1. 修复WebSocket连接的异步初始化

将_setup_streamer_websocket改为异步方法,所有websockets操作都使用await:

@property
async def streamer(self) -> ClientConnection:
    if not self._streamer:
        self._streamer = await self._setup_streamer_websocket()
    return self._streamer

async def _setup_streamer_websocket(self) -> ClientConnection:
    websocket = await connect(self.quote_streamer_url)

    setup_message_payload = {
        "type": "SETUP",
        "channel": 0,
        "keepaliveTimeout": 60,
        "acceptKeepaliveTimeout": 60,
        "version": "0.1",
    }

    await websocket.send(dumps(setup_message_payload))

    auth_message_payload = {
        "type": "AUTH",
        "channel": 0,
        "token": self.quote_streamer_token,
    }

    await websocket.send(dumps(auth_message_payload))

    return websocket

2. 修正消息处理循环的异步调用

async def message_handler(self):
    while True:
        message = await self.streamer.recv()
        await self._process_message(message)

3. 调整主程序初始化逻辑

确保先建立WebSocket连接再启动消息处理器:

async def main():
    api = ApiClient(
        uid=_uid,
        remember_token=_remember_token
    )
    # 先初始化WebSocket连接
    await api.streamer
    # 启动消息处理任务
    asyncio.create_task(api.message_handler())

    await api.request_channel(1)
    event = api.open_channels.get(channel)
    await event.wait()
    log.info("Subscribing options quotes for tracked symbols.")
    # 后续代码保持不变

4. 修复request_channel中的send操作

将request_channel改为异步方法,因为streamer.send()是异步操作:

async def request_channel(self, channel: int):
    if self.open_channels.get(channel):
        log.info(f"Channel {channel} was already requested. Skipping request.")
        return

    self.open_channels[channel] = asyncio.Event()
    await self.streamer.send(
        dumps(
            {
                "type": "CHANNEL_REQUEST",
                "channel": channel,
                "service": "FEED",
                "parameters": {"contract": "AUTO"},
            }
        )
    )
    log.debug(f"Requested opening of channel {channel}")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 04:00:55