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

如何优雅刷新aiohttp Client Session并解决未关闭会话警告?

解决aiohttp令牌刷新时的未关闭会话警告

问题根源

你遇到的unclosed client session和unclosed connector警告,本质是刷新令牌时直接创建新的ClientSession,但旧会话未被显式关闭。aiohttp的会话持有连接池、TCP连接器等资源,必须通过await session.close()主动释放,否则残留的未关闭资源会触发警告。

核心修正方案

  1. 异步处理会话生命周期:关闭会话是异步操作,必须用await执行;刷新令牌的POST请求也需用异步方式完成,避免阻塞事件循环。
  2. 集中管理状态:将会话、令牌、登录时间等状态放到类属性中统一维护,避免在请求函数里临时替换会话导致资源泄漏。
  3. 优雅重试逻辑:刷新会话后重试当前请求,同时覆盖定时检测和请求返回401两种令牌失效场景。

修正后的完整代码示例

import time
import asyncio
import aiohttp

class Request:
    def __init__(self, url: str, method: str="get", payload: dict=None, headers: dict=None):
        self.url: str = url
        self.method: str = method
        self.payload: dict = payload or dict()
        self.headers: dict = headers or dict()

class Response:
    def __init__(self, url: str, status: int, payload: dict=None, error: bool=False, text: str=None):
        self.url: str = url
        self.status: int = status
        self.payload: dict = payload or dict()
        self.error: bool = error
        self.text: str = text or ''

class APIClient:
    def __init__(self, initial_tokens: dict):
        self.tokens = initial_tokens
        self.login_time = time.time()
        self.client_session = None
        # 异步初始化会话
        asyncio.create_task(self._init_session())

    async def _init_session(self):
        """初始化带认证头的客户端会话"""
        headers = {
            'Content-Type': 'application/json',
            'Authorization': f"Bearer {self.tokens['accessToken']}"
        }
        self.client_session = aiohttp.ClientSession(headers=headers)

    async def _refresh_tokens(self) -> dict:
        """异步请求刷新令牌,使用临时会话避免依赖失效的原有会话"""
        refresh_url = "https://api.example.com/refresh-token"  # 替换为实际刷新地址
        refresh_payload = {
            'accessToken': self.tokens['accessToken'],
            'refreshToken': self.tokens['refreshToken']
        }
        async with aiohttp.ClientSession() as temp_session:
            async with temp_session.post(refresh_url, json=refresh_payload) as resp:
                resp.raise_for_status()
                return await resp.json()

    async def _refresh_session(self):
        """关闭旧会话并创建带新令牌的新会话"""
        # 先关闭旧会话释放资源
        if self.client_session:
            await self.client_session.close()
        # 获取新令牌并更新状态
        self.tokens = await self._refresh_tokens()
        self.login_time = time.time()
        # 创建新会话
        headers = {
            'Content-Type': 'application/json',
            'Authorization': f"Bearer {self.tokens['accessToken']}"
        }
        self.client_session = aiohttp.ClientSession(headers=headers)

    async def fetch(self, request: Request) -> Response:
        """核心请求方法,处理令牌刷新和重试"""
        # 定时检测令牌是否过期
        if time.time() >= self.login_time + 10:
            print("Refreshing tokens by time check...")
            await self._refresh_session()

        try:
            # 根据请求方法选择参数传递方式
            request_kwargs = {}
            if request.method.lower() == 'get':
                request_kwargs['params'] = request.payload
            else:
                request_kwargs['json'] = request.payload
            
            async with self.client_session.request(
                request.method, 
                request.url,
                headers=request.headers,
                **request_kwargs
            ) as response:
                text = await response.text(encoding='utf-8')
                # 处理非JSON响应
                try:
                    content = await response.json()
                except aiohttp.ContentTypeError:
                    content = {}
                response.raise_for_status()
                return Response(request.url, response.status, payload=content, text=text)
        except aiohttp.ClientResponseError as e:
            # 处理401令牌失效,强制刷新并重试
            if e.status == 401:
                print("Token expired, refreshing by 401 response...")
                await self._refresh_session()
                return await self.fetch(request)
            # 其他HTTP错误处理
            return Response(request.url, e.status, error=True, text=str(e))
        except Exception as e:
            # 通用异常处理
            return Response(request.url, 500, error=True, text=str(e))

# 使用示例
async def main():
    initial_tokens = {
        'accessToken': 'your_initial_access_token',
        'refreshToken': 'your_initial_refresh_token'
    }
    client = APIClient(initial_tokens)
    # 等待会话初始化完成
    await asyncio.sleep(0.1)
    
    # 发起请求
    req = Request(url="https://api.example.com/data", method="get", payload={"key": "value"})
    resp = await client.fetch(req)
    print(resp.payload)
    
    # 程序结束时主动关闭会话
    await client.client_session.close()

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

关键细节说明

  • 资源释放优先级:刷新会话时先关闭旧会话,再创建新会话,确保连接池等资源被及时回收。
  • 双场景失效处理:既通过定时检测提前刷新令牌,又处理请求返回的401错误,覆盖主动过期和被动失效两种情况。
  • 临时会话刷新令牌:刷新令牌时使用独立的临时会话,避免原有会话已失效导致刷新请求失败。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 04:24:25