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

Python Async处理大量Zip文件时内存泄漏(OOM)原因排查

内存累积导致OOM的问题分析与解决

问题根源

你的代码存在几个核心内存泄漏/累积点:

  • 全量加载Zip文件到内存:每次把整个Blob Zip文件读入io.BytesIO(),如果Zip体积大,并行任务叠加后内存直接暴涨。
  • asyncio.gather缓存所有结果:asyncio.gather会等待所有任务完成并收集所有返回的DataFrame,数百万条文件名的DataFrame全部堆在内存里,必然触发OOM。
  • 重复创建Azure客户端:每个zip_reader都重新创建ClientSecretCredential和BlobServiceClient,这些客户端的内部资源未及时释放,造成内存泄漏。
  • 无效的内存清理:cleanup里的gc.collect()时机不对,且未显式解除大对象的引用,垃圾回收无法有效工作。

优化方案

1. 复用Azure客户端资源

在类初始化时创建一次Credential和BlobServiceClient,避免重复创建:

def __init__(self):
    self.credential = ClientSecretCredential(TENANT, CLIENTID, CLIENTSECRET)
    self.blob_service_client = BlobServiceClient(
        account_url="https://blob1.blob.core.windows.net/",
        credential=self.credential,
        max_single_get_size=64 * 1024 * 1024,
        max_chunk_get_size=32 * 1024 * 1024
    )

2. 流式处理Zip,避免全量加载

用临时文件替代内存中的BytesIO存储下载的Zip内容,处理完立即删除:

import tempfile
import os
import gc

async def zip_reader(self, blobFileName, blobEndPoint, semaphore):
    try:
        async with semaphore:
            logger.info(f"Starting: {blobFileName}, {blobEndPoint}")
            async with self.blob_service_client.get_blob_client(container=blobEndPoint, blob=blobFileName) as blob_client:
                # 用临时文件存储Zip,减少内存占用
                with tempfile.NamedTemporaryFile(delete=False) as temp_file:
                    stream = await blob_client.download_blob(max_concurrency=10)  # 降低并发数,减少内存峰值
                    await stream.readinto(temp_file)
                    temp_file_path = temp_file.name
                
                # 读取Zip内文件名后立即关闭Zip
                with ZipFile(temp_file_path, 'r') as f:
                    file_list = f.namelist()
                
                # 立即删除临时文件
                os.unlink(temp_file_path)

                # 若必须用DataFrame,考虑直接写入外部存储而非留存内存
                t_df = pd.DataFrame({'fileList': file_list})
                t_df['blobFileName'] = blobFileName
                t_df['blobEndPoint'] = blobEndPoint

                logger.info(f"Completed: {blobFileName}")

                # 显式解除大对象引用,触发垃圾回收
                del file_list
                gc.collect()

                return t_df
    except Exception as e:
        logger.error(f"Error processing {blobFileName}: {str(e)}")
        # 异常时确保临时文件被清理
        if 'temp_file_path' in locals():
            os.unlink(temp_file_path)
        raise

3. 避免一次性缓存所有结果

用asyncio.as_completed替代asyncio.gather,处理完一个任务就输出/写入外部存储,释放内存:

async def async_file_as_bytes_generator(self, blobFileName, blobEndPoint, semaphore_value):
    semaphore = asyncio.Semaphore(value=semaphore_value)
    tasks = [self.zip_reader(fn, ep, semaphore) for fn, ep in zip(blobFileName, blobEndPoint)]
    
    # 逐个处理完成的任务,避免内存堆积
    for task in asyncio.as_completed(tasks):
        t_df = await task
        # 这里直接将结果写入外部存储(如数据库/CSV),然后立即删除DataFrame
        # 示例:t_df.to_csv('output.csv', mode='a', header=False, index=False)
        del t_df
        gc.collect()

4. 其他关键优化

  • 降低max_concurrency:当前设置为25,过高的并发会导致多个大文件同时下载,建议根据服务器内存调整到5-10。
  • 取消不必要的DataFrame:如果只是需要处理文件名,直接遍历f.namelist()逐个处理,完全不需要用Pandas存储,能大幅减少内存占用。

内容的提问来源于stack exchange,提问作者John Stud

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 17:15:43