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

基于asyncio、aiohttp、boto3的内存泄漏排查求助

内存泄漏问题分析与解决方案

可能的泄漏根源

  1. aiohttp资源未正确释放:未通过上下文管理ClientSession或请求响应,导致连接、请求对象残留,连带持有Traceback等关联对象。
  2. boto3客户端管理不当:同步boto3在异步环境中频繁创建实例且未关闭,或未使用异步客户端(如aioboto3),导致客户端对象累积。
  3. Traceback/异常对象被持久化引用:错误处理逻辑中将完整异常(含Traceback)存入全局容器(如错误列表),或日志配置保留过多追踪信息,导致对象无法被GC回收。
  4. 批量处理数据残留:批量任务的中间结果、任务对象被全局变量或外层作用域容器持有,未及时清空,形成内存堆积。

排查方向

  • 检查aiohttp代码:确认ClientSession是否全局复用并最终关闭,所有请求是否通过async with session.get(url)上下文管理响应对象。
  • 检查boto3调用:若用同步boto3,是否复用客户端实例而非每次请求新建;若用aioboto3,是否通过async with管理客户端生命周期。
  • 检查错误处理逻辑:是否存在将sys.exc_info()返回的Traceback、完整异常对象存入全局集合的情况,是否仅保留必要错误信息。
  • 检查批量数据容器:处理完每批任务后,是否清空存储中间结果的列表、字典,避免容器无限增长。

可行解决方案

1. 规范aiohttp资源管理

  • 全局复用单个ClientSession实例,程序退出时显式关闭:
    async def main():
        async with aiohttp.ClientSession() as session:
            while True:
                urls = get_latest_doc_urls()
                await process_batch(urls, session)
    
  • 所有请求使用上下文管理响应对象,确保资源自动释放:
    async def fetch_doc(session, url):
        async with session.get(url) as resp:
            return await resp.read()
    

2. 优化boto3异步调用

  • 替换为aioboto3异步客户端,通过上下文管理确保资源释放:
    async def upload_to_s3(content, key):
        async with aioboto3.client('s3') as s3:
            await s3.put_object(Bucket='my-bucket', Key=key, Body=content)
    
  • 若必须使用同步boto3,全局复用客户端实例,避免频繁创建:
    s3_client = boto3.client('s3')
    
    async def upload_to_s3(content, key):
        await asyncio.to_thread(
            s3_client.put_object,
            Bucket='my-bucket', Key=key, Body=content
        )
    

3. 清理Traceback与异常引用

  • 避免存储完整异常对象,仅保留错误描述、代码等必要信息:
    # 错误示例:保留完整异常
    errors.append(exc)
    # 正确做法:仅存储关键信息
    errors.append({"url": url, "error": str(exc), "code": exc.status})
    
  • 若使用traceback.format_exc()生成错误日志,确保日志系统不会无限缓存旧日志条目。

4. 主动清理批量数据

  • 每批任务处理完成后,显式清空中间结果容器:
    async def process_batch(urls, session):
        tasks = [fetch_doc(session, url) for url in urls]
        contents = await asyncio.gather(*tasks)
        for idx, content in enumerate(contents):
            converted = convert_doc(content)
            await upload_to_s3(converted, f'doc_{idx}.txt')
        # 清空临时结果,断开引用
        del contents[:]
    

5. 定位循环引用

  • 启用GC调试模式,查看无法回收的对象及引用链:
    import gc
    gc.set_debug(gc.DEBUG_SAVEALL)
    # 运行程序一段时间后
    unreachable = gc.collect()
    print(f"Unreachable objects: {unreachable}")
    for obj in gc.garbage:
        print(type(obj), obj)
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 08:48:11