复用aiohttp ClientSession触发RuntimeError:Timeout上下文需在任务内
问题:复用aiohttp.ClientSession触发RuntimeError:Timeout context manager should be used inside a task
尝试复用aiohttp.ClientSession发起多个请求时触发RuntimeError: Timeout context manager should be used inside a task,但每次请求新建会话可正常运行。使用环境:Python 3.10.2、aiohttp 3.8.5。
错误回溯
RuntimeError: Timeout context manager should be used inside a task Unclosed client session プロセスは終了コード 1 で終了しました
原代码
import asyncio from abc import ABC, abstractmethod import aiohttp from typing import Dict, List, Any def get_headers(extra: Dict[str, Any] = {}) -> Dict[str, str]: headers = { "accept": "text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,image/apng,*/*;q=0.8,application/signed-exchange;v=b3;q=0.9", "accept-language": "en-GB,en;q=0.9,ja-JP;q=0.8,ja;q=0.7,en-US;q=0.6", "user-agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/104.0.5112.102 Safari/537.36", } for key, val in extra.items(): headers[key] = val return headers class Scraper(ABC): session: aiohttp.ClientSession = None def __init__(self): if not self.session: self.session = asyncio.get_event_loop().run_until_complete(self.get_session()) @classmethod async def get_session(cls): cls.session = aiohttp.ClientSession() return cls.session async def __aenter__(self): return self async def __aexit__(self, exc_type, exc_value, traceback): if self.session: await self.session.close() class Animepahe(Scraper): _SITE_NAME: str = "animepahe" site_url: str = "https://animepahe.ru" api_url: str = "https://animepahe.ru/api" manifest_header = get_headers({"referer": "https://kwik.cx", "origin": "https://kwik.cx"}) async def get_api(self, data: dict, headers: dict = get_headers()) -> dict: resp = await self.get(self.api_url, data, headers) return await resp.json() async def get(self, url: str, data=None, headers: dict = get_headers()) -> aiohttp.ClientResponse: data = {} or data err, tries = None, 0 while tries < 10: try: async with self.session.get(url=url, params=data, headers=headers) as resp: if resp.status != 200: err = f"request failed with status: {resp.status}\n err msg: {resp.content}" logging.error(f"{err}\nRetrying...") raise aiohttp.ClientResponseError(None, None, message=err) return resp except (aiohttp.ClientOSError, asyncio.TimeoutError, aiohttp.ServerDisconnectedError, aiohttp.ServerTimeoutError): await asyncio.sleep(choice([5, 4, 3, 2, 1])) # randomly await tries += 1 continue raise aiohttp.ClientResponseError(None, None, message=err) if __name__ == "__main__": Scraper() async def main(): scraper = Animepahe() print(await scraper.get_api({"m": "search", "q": "attack"})) asyncio.run(main())
问题原因
- 同步上下文创建Session:
Scraper的__init__方法通过asyncio.get_event_loop().run_until_complete在同步环境中创建ClientSession,而aiohttp的超时管理器依赖异步任务上下文,后续在异步函数中使用该Session时会触发上下文冲突。 - 类属性Session共享:
session是类属性,会被所有子类实例共享,易引发会话复用或关闭冲突。 - 可变默认参数污染:
get_headers的默认参数extra={}是可变对象,多次调用会共享同一个字典,导致意外的头部污染。 - 缺失依赖导入:代码未导入
logging和random.choice,运行时会报错。 - 提前实例化Scraper:
if __name__ == "__main__"中提前实例化Scraper(),在asyncio.run之前同步创建Session,加剧上下文冲突。
修复方案
修复后完整代码
import asyncio import logging from random import choice from abc import ABC, abstractmethod import aiohttp from typing import Dict, List, Any, Optional def get_headers(extra: Optional[Dict[str, Any]] = None) -> Dict[str, str]: headers = { "accept": "text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,image/apng,*/*;q=0.8,application/signed-exchange;v=b3;q=0.9", "accept-language": "en-GB,en;q=0.9,ja-JP;q=0.8,ja;q=0.7,en-US;q=0.6", "user-agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/104.0.5112.102 Safari/537.36", } if extra: headers.update(extra) return headers class Scraper(ABC): def __init__(self): self.session: Optional[aiohttp.ClientSession] = None async def init_session(self): if not self.session: self.session = aiohttp.ClientSession() return self.session async def __aenter__(self): await self.init_session() return self async def __aexit__(self, exc_type, exc_value, traceback): if self.session: await self.session.close() class Animepahe(Scraper): _SITE_NAME: str = "animepahe" site_url: str = "https://animepahe.ru" api_url: str = "https://animepahe.ru/api" manifest_header = get_headers({"referer": "https://kwik.cx", "origin": "https://kwik.cx"}) async def get_api(self, data: dict, headers: Optional[dict] = None) -> dict: headers = headers or get_headers() resp = await self.get(self.api_url, data, headers) return await resp.json() async def get(self, url: str, data: Optional[dict] = None, headers: Optional[dict] = None) -> aiohttp.ClientResponse: data = data or {} headers = headers or get_headers() err, tries = None, 0 while tries < 10: try: async with self.session.get(url=url, params=data, headers=headers) as resp: if resp.status != 200: err = f"request failed with status: {resp.status}\n err msg: {await resp.text()}" logging.error(f"{err}\nRetrying...") raise aiohttp.ClientResponseError(None, None, message=err) return resp except (aiohttp.ClientOSError, asyncio.TimeoutError, aiohttp.ServerDisconnectedError, aiohttp.ServerTimeoutError): await asyncio.sleep(choice([5, 4, 3, 2, 1])) tries += 1 continue raise aiohttp.ClientResponseError(None, None, message=err) if __name__ == "__main__": async def main(): async with Animepahe() as scraper: print(await scraper.get_api({"m": "search", "q": "attack"})) asyncio.run(main())
关键修复点
- 异步上下文管理Session:使用
async with初始化和关闭Session,确保全程在异步上下文内操作。 - 实例化Session:将
session改为实例属性,每个Animepahe实例拥有独立会话,避免跨实例冲突。 - 修复可变默认参数:将
get_headers的默认参数改为None,内部初始化空字典,避免状态污染。 - 补全依赖导入:添加
logging和random.choice的导入。 - 移除提前实例化:删除
if __name__ == "__main__"中的Scraper()调用,避免同步创建Session。 - 修复响应内容获取:将
resp.content改为await resp.text(),正确获取响应文本内容。
内容的提问来源于stack exchange,提问作者CosmicOppai
相关产品推荐
相关产品推荐

