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
问题根源分析
- WebSocket异步API被同步调用:websockets库的
connect()、send()、recv()都是异步方法,当前代码直接同步调用这些方法,得到的是协程对象而非实际连接/结果,导致消息收发逻辑完全失效。 - 消息处理循环的异步调用错误:
message_handler中直接传递self.streamer.recv()给_process_message,但recv()是异步方法,必须用await获取结果后再传入。 - 事件循环调度异常:错误的同步操作打乱了事件循环的任务调度,即使
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
相关产品推荐
相关产品推荐

