如何优雅刷新aiohttp Client Session并解决未关闭会话警告?
解决aiohttp令牌刷新时的未关闭会话警告
问题根源
你遇到的unclosed client session和unclosed connector警告,本质是刷新令牌时直接创建新的ClientSession,但旧会话未被显式关闭。aiohttp的会话持有连接池、TCP连接器等资源,必须通过await session.close()主动释放,否则残留的未关闭资源会触发警告。
核心修正方案
- 异步处理会话生命周期:关闭会话是异步操作,必须用
await执行;刷新令牌的POST请求也需用异步方式完成,避免阻塞事件循环。 - 集中管理状态:将会话、令牌、登录时间等状态放到类属性中统一维护,避免在请求函数里临时替换会话导致资源泄漏。
- 优雅重试逻辑:刷新会话后重试当前请求,同时覆盖定时检测和请求返回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
相关产品推荐
相关产品推荐

