Asyncio+Aiohttp执行超50万异步请求时内存泄漏问题求助
解决Aiohttp+Asyncio批量请求的内存泄漏问题
你的代码存在多个导致内存持续增长的问题,以下是具体问题分析和修复方案:
问题分析
- 重复创建ClientSession:
fetch函数明明已经接收了外部传入的session,却每次请求都新建一个ClientSession,且没有通过async with正确管理生命周期,导致大量未关闭的会话对象堆积在内存中。 - 重复调用
asyncio.gather:fetch_all中先await asyncio.gather(*tasks)完成所有任务,之后又再次调用await asyncio.gather(*tasks),这不仅是冗余操作,还会让tasks列表一直持有已完成任务的引用,无法被垃圾回收器回收,持续占用内存。 - Semaphore使用方式错误:在循环中对每个URL都执行
async with sem,这会导致任务串行创建,完全没利用并发优势,还可能引发任务对象堆积。 - 重复读取响应内容:同时调用
resp.text()和resp.read()会把响应内容在内存中存储两份(二进制字节流和字符串),大量请求下会额外消耗大量内存。
修复后的代码
import asyncio from asyncio import Semaphore import aiohttp from bs4 import BeautifulSoup async def fetch(session, sem, url): try: headers = {} # 替换为你的实际headers params = {} # 替换为你的实际params async with sem: async with session.get(url, headers=headers, params=params, verify_ssl=False) as resp: # 只读取文本内容,不需要同时存二进制和字符串 text = await resp.text() return text except Exception as e: print(f"url: {url} error happened: {str(e)}") # 异常时返回None,避免后续处理报错 return None async def fetch_all(urls): # 限制并发数 sem = Semaphore(100) # 使用TCPConnector控制连接池,避免连接泄漏 connector = aiohttp.TCPConnector(limit=100, force_close=True) async with aiohttp.ClientSession(connector=connector, cookie_jar=aiohttp.DummyCookieJar()) as session: tasks = [] for url in urls: task = asyncio.create_task(fetch(session, sem, url)) tasks.append(task) # 只执行一次gather,获取所有结果 datas = await asyncio.gather(*tasks, return_exceptions=True) return datas def get_result(data, chupindict): if not data: return try: soup = BeautifulSoup(data, 'html.parser') name = soup.find('div', class_='name').text.strip() chupin = soup.find('div', class_="panel-wrapper", id="出品").text.strip() chupindict[name] = chupin except Exception as e: print(f"解析数据出错: {str(e)}") if __name__ == "__main__": urls = [] chupindict = {} # 生成测试URL(替换为你的实际URL生成逻辑) for i in range(0, 500000): url = f"http://example.com/test/{i}" urls.append(url) loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) try: datas = loop.run_until_complete(fetch_all(urls)) # 逐个处理结果,避免一次性持有所有数据 for data in datas: get_result(data, chupindict) finally: loop.close()
关键修改说明
- 复用ClientSession:全局只创建一个
ClientSession,所有请求复用这个会话,通过TCPConnector控制连接池大小,避免连接泄漏。 - 正确使用Semaphore:将信号量传入
fetch函数,在请求内部获取信号量,确保并发数被正确控制,同时任务可以批量创建。 - 单次调用gather:只执行一次
asyncio.gather获取所有结果,避免任务引用堆积。 - 避免重复读取响应:只保留需要的响应格式(这里用文本),减少内存占用。
- 结果分批处理:逐个处理返回的结果,而不是一次性持有所有响应数据,降低峰值内存占用。
- 异常处理优化:在
fetch和get_result中都增加了异常捕获,避免单个请求/解析失败导致程序崩溃,同时减少无效内存占用。
内容的提问来源于stack exchange,提问作者Max Chen
相关产品推荐
相关产品推荐

