使用asyncio爬取数据时内存利用率持续升高问题排查
问题根因
- 任务引用无限累积:
tasks列表定义在header循环外层,第一轮header对应的任务执行完成后,任务对象引用一直残留在列表中不会被回收;外层while True无限循环每跑一轮,都会往同一个列表里重复追加新任务,内存占用随运行时间线性上涨。 - 同步Redis客户端在异步环境中使用不当:代码中使用同步版
redis.Redis客户端,在协程内直接调用同步方法会阻塞整个事件循环,且同步客户端的连接、缓冲区不受asyncio事件循环管理,长期运行会出现连接泄漏、缓冲内存无法释放的问题。代码中已导入aioredis但未实际使用。 - ClientSession生命周期管理错误:在header循环内部通过
async with创建ClientSession,但会话创建后就往任务列表追加了引用session的协程,async with代码块退出时session会被立即关闭,此时未执行完的协程会持有已失效的session资源,相关异常被裸except吞掉后,未关闭的连接、响应残留在内存中无法回收。 - 资源释放逻辑缺失+裸异常捕获:
- aiohttp请求未使用上下文管理器管理响应对象,请求超时、报错时未显式关闭响应,连接无法归还到连接池,持续占用内存。
- 所有异常分支都用
except: pass直接吞掉错误,既不打日志也不释放资源,泄漏的资源会持续累积。 get函数的重试逻辑有缺陷,第一次请求抛出异常时,未清理残留连接就直接发起第二次请求,进一步加剧连接泄漏。
- 业务逻辑错误导致异常频发:
publish函数中对字符串类型的入参调用data.empty(该属性是pandas对象特有),判断逻辑恒定抛出异常,Redis写入逻辑实际从未正常执行,异常被吞的同时也会残留部分调用栈内存。 - 调试模式长期开启:
asyncio.run传入了debug=True参数,asyncio调试模式会保留所有任务的执行栈、上下文快照用于调试,长期运行下这部分内存会持续增长不会自动释放。 - 无用全局变量驻留:代码开头定义的
RESULTS = []、result_dict = {}全程未被使用,作为全局变量永久驻留内存,若后续误写入数据会直接加剧内存泄漏。
修复优化方案
- 调整任务列表作用域:将
tasks = []的定义移到单个ClientSession对应的循环内部,每一批任务执行完成后,tasks列表随作用域结束自动回收,避免跨批次、跨循环累积任务引用。 - 替换为异步Redis客户端:删除同步Redis客户端初始化代码,使用已导入的
aioredis创建异步连接,连接全局初始化一次在所有协程间复用,所有Redis操作都用await调用,避免事件循环阻塞和连接泄漏,初始化示例:
# 替换原同步Redis初始化代码 async def init_redis(): return await aioredis.from_url( "redis://你的redis地址:6379", username="default", password="你的redis密码", decode_responses=True, max_connections=20 # 配置连接数上限,避免连接泄漏 )
- 修正ClientSession生命周期:将ClientSession的
async with块包裹住对应批次的任务创建、gather执行全流程,保证所有持有session引用的任务执行完成后,再关闭会话。创建Session时配置连接池上限,和信号量数值对齐:
connector = aiohttp.TCPConnector(limit=20) # 和信号量并发数一致 async with ClientSession(headers=value, connector=connector) as session: tasks = [] # tasks移到这里,每批任务单独创建列表 for i in range(0, 5): URLS = f"https://pokeapi.co/api/v2/pokemon/{i}" tasks.append(asyncio.create_task(run_program(URLS, session, semaphore, redis_conn))) await asyncio.gather(*tasks) # 会话退出后tasks列表自动回收
- 补全资源释放逻辑,禁止裸异常捕获:
- 所有aiohttp请求使用
async with上下文管理器管理响应,保证响应无论成功失败都会被正确关闭,连接归还连接池,get函数修正为:
- 所有aiohttp请求使用
async def get(url, session): try: async with session.request(method="GET", url=url, timeout=1) as response: response.raise_for_status() # 主动抛HTTP错误状态码异常 pokemon = await response.json() return pokemon["name"] except Exception as err: # 第一次请求失败重试,上下文管理器已经自动释放第一次的连接 async with session.request(method="GET", url=url, timeout=3) as response: response.raise_for_status() pokemon = await response.json() return pokemon["name"]
- 所有
except块明确捕获异常类型,禁止无差别吞所有异常,异常分支打印日志,显式释放占用的资源。 - 修正publish函数逻辑:删除错误的
data.empty判断,直接判断入参非空即可,Redis操作用异步客户端加await:
async def publish(data, redis_conn): if data: try: keyName = "channelName" await redis_conn.set(keyName, data) except Exception as e: logger.error(f"写入Redis失败: {str(e)}")
- 关闭生产环境调试模式:
asyncio.run调用时去掉debug=True参数,或者显式传debug=False,避免调试信息长期驻留内存。 - 清理无效代码:删除未使用的全局变量
RESULTS、result_dict,删除未用到的导入(如multiprocessing.Semaphore、requests.exceptions.HTTPError等),减少不必要的内存驻留。 - 增加定期GC回收:在
while True循环每轮执行完成后,调用gc.collect()手动触发一次垃圾回收,及时释放循环引用的内存。
服务器内存变化趋势参考

内容的提问来源于stack exchange,提问作者Yadvendar Singh
相关产品推荐
相关产品推荐

