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
相关产品推荐
相关产品推荐

