Asyncio.ws_connect未收到连接关闭消息问题排查
我编写了如下基于aiohttp的WebSocket客户端代码:
async def _connect_session(self): websocket_url, headers, cookies = self._get_connect_data() try: async with aiohttp.ClientSession(cookies=cookies) as session: logger.debug( f"Connecting to WebSocket using proxy {self._api.get_proxy().url}" ) async with session.ws_connect( websocket_url, headers=headers, proxy=self._api.get_proxy().url, verify_ssl=False, ) as connection: self.status = True self._connection = connection self._wait_connect_task.cancel() self._matches_controller_task = asyncio.create_task(self._matches_controller()) logger.debug("Connected to websocket") self.ping_task = asyncio.create_task(self._ping(connection)) await self._listen(connection) logger.error(f"WsApp closed {self}") self._destruct(WsClosed()) except Exception as e: logger.exception(f"Exception in WsApp.connect {self}") self._destruct(e) async def _ping(self, connection: ClientWebSocketResponse): while True: try: await asyncio.sleep(19) await connection.send_str("2") except Exception as e: self._destruct(e) async def _listen(self, connection: ClientWebSocketResponse): while True: msg = await connection.receive() if msg.type == WSMsgType.TEXT: data: str = msg.data self._message_processor(data) elif msg.type in ( WSMsgType.CLOSE, WSMsgType.CLOSED, WSMsgType.CLOSING, WSMsgType.ERROR, ): raise WsClosed() def _destruct(self, err: BaseException): self._callbacks.error_callback(self, err) if self._run_task is not None: self._run_task.cancel() if self._matches_controller_task is not None: self._matches_controller_task.cancel() if self.ping_task is not None: self.ping_task.cancel()
当尝试发送ping消息时连接频繁断开,但_listen方法始终无法收到连接关闭相关消息,日志中仅出现_ping方法的如下异常:
Traceback (most recent call last): File "/opt/adapter/tipsport/tipsport_adapter/pipeline/single_match/data_providers/ws_provider/ws_app.py", line 139, in _ping await connection.send_str("2") File "/opt/adapter/.venv/lib/python3.11/site-packages/aiohttp/client_ws.py", line 225, in send_str await self._writer.send(data, binary=False, compress=compress) File "/opt/adapter/.venv/lib/python3.11/site-packages/aiohttp/http_websocket.py", line 726, in send await self._send_frame(message, WSMsgType.TEXT, compress) File "/opt/adapter/.venv/lib/python3.11/site-packages/aiohttp/http_websocket.py", line 626, in _send_frame raise ConnectionResetError("Cannot write to closing transport") ConnectionResetError: Cannot write to closing transport
无任何来自_listen的消息或错误,请问为何会出现这种情况?
核心原因
TCP连接被动断开,无WebSocket关闭帧:当底层TCP连接被对方直接切断(比如超时、网络波动、服务器强制终止),不会发送标准的WebSocket关闭帧。此时
connection.receive()会一直阻塞,因为没有收到任何WebSocket层面的消息,自然触发不了你监听的CLOSE/CLOSED等类型。只有当发送ping时,操作系统检测到连接已断,才会抛出ConnectionResetError。_listen任务无限阻塞:_listen里的await connection.receive()是无等待时长的阻塞操作,如果没有消息推送,它会一直挂起。TCP被动断开时没有WebSocket层面的信号,所以_listen不会输出任何内容或触发异常。_destruct未处理_listen任务:在_destruct中你取消了多个任务,但没有处理_listen对应的任务。当ping_task抛出异常触发_destruct后,_listen仍会阻塞,直到ClientSession或ws_connect的上下文管理器退出才会终止,不会走你_listen里的异常逻辑。
修复方案
给
_listen添加超时检查:用asyncio.wait_for包装receive(),定期检查连接状态,避免无限阻塞:async def _listen(self, connection: ClientWebSocketResponse): while True: try: # 设置10秒超时,定期检查连接 msg = await asyncio.wait_for(connection.receive(), timeout=10) if msg.type == WSMsgType.TEXT: data: str = msg.data self._message_processor(data) elif msg.type in ( WSMsgType.CLOSE, WSMsgType.CLOSED, WSMsgType.CLOSING, WSMsgType.ERROR, ): raise WsClosed() except asyncio.TimeoutError: # 超时后检查连接是否已关闭 if connection.closed: raise WsClosed() # 未关闭则继续等待消息 continue在
_destruct中取消_listen任务:先在_connect_session中保存_listen任务引用,再在_destruct中取消:# 在_connect_session中保存listen任务 self._listen_task = asyncio.create_task(self._listen(connection)) await self._listen_task # 在_destruct方法中添加 if self._listen_task is not None: self._listen_task.cancel()改用aiohttp内置心跳机制:
ws_connect支持heartbeat参数,可自动发送ping帧,比手动实现更可靠:async with session.ws_connect( websocket_url, headers=headers, proxy=self._api.get_proxy().url, verify_ssl=False, heartbeat=19, # 自动每19秒发送ping ) as connection:
内容的提问来源于stack exchange,提问作者Vladimir Vasiliev

