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

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的消息或错误,请问为何会出现这种情况?

分析与解决

核心原因

  1. TCP连接被动断开,无WebSocket关闭帧:当底层TCP连接被对方直接切断(比如超时、网络波动、服务器强制终止),不会发送标准的WebSocket关闭帧。此时connection.receive()会一直阻塞,因为没有收到任何WebSocket层面的消息,自然触发不了你监听的CLOSE/CLOSED等类型。只有当发送ping时,操作系统检测到连接已断,才会抛出ConnectionResetError。

  2. _listen任务无限阻塞:_listen里的await connection.receive()是无等待时长的阻塞操作,如果没有消息推送,它会一直挂起。TCP被动断开时没有WebSocket层面的信号,所以_listen不会输出任何内容或触发异常。

  3. _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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 08:20:21