如何强制释放Python进程内存以避免累积分配引发MemoryError?
环境信息
- 操作系统:RHEL 8.9
- 可用内存:128GB
- Python版本:3.11.5
问题描述
加载3-5GB、含3500万+行的CSV文件(每行是3-10个整数的列表)及若干JSON文件到DataFrame,其中CSV对应big_df。加载与基础预处理完成后,主进程内存占用(VmRSS)约7.9GB。
通过multiprocessing.Pool(processes=3, maxtasksperchild=5)启动3个工作进程,调用pool.map分发任务:每个进程处理big_df的1000行数据,通过嵌套循环(1000-5000次迭代)做集合交集判断,将符合条件的结果存入列表后返回主进程。
执行1200-2500个任务后,工作进程在序列化结果列表返回主进程时触发MemoryError。日志显示:进程初始内存均为7.9GB,触发错误时3个进程的VmRSS分别达47.8GB、28GB、26.5GB,内存持续上升且未释放给操作系统。
已尝试的无效方案
- 减少进程数(从6降至3/2):仅延迟错误触发时间,未解决内存持续增长问题
- 将
big_df设为全局变量:依赖Linux写时复制(copy-on-write)避免重复复制,但内存占用仍持续上升 - 设置
maxtasksperchild=5:期望进程完成5次任务后重置释放内存,无效
伪代码
主进程
global big_df, smaller_df # 分别从CSV/JSON加载 chunk_indices = [(start_idx, start_idx+1000) for start_idx in range(0, len(big_df), 1000)] with multiprocessing.Pool(processes=3, maxtasksperchild=5) as p: result_lists = p.map(worker, chunk_indices) p.close() p.join() # 从result_lists创建DataFrame并写入输出CSV
工作进程
def worker(args): logging.basicConfig(filename="logfile.log", level=logging.DEBUG) start_idx, end_idx = args result = [] for i in range(start_idx, end_idx): list_i = big_df[i] for j in range(len(smaller_df)): if len(list_i & smaller_df[j]): result.append(smaller_df[j]) logging.debug("time: %s, sizeof results: %s", ctime(), result.__sizeof__()) return result
日志显示:工作进程的result数组最大可达9GB,通常为3-5GB,若任务完成后内存能正常释放则不会触发错误。
解决方案建议
1. 避免大列表跨进程传递,改用文件直接写入
每次任务返回的大列表是内存堆积的核心原因,可改为工作进程直接将结果写入临时文件,主进程最后合并文件:
- 工作进程处理完chunk后,将结果写入独立临时文件(如
result_{start_idx}.csv) - 主进程在所有任务完成后,遍历临时文件合并为最终输出CSV
- 此方式避免了大列表序列化/反序列化的开销,且工作进程写完即可清空
result释放内存
2. 手动触发垃圾回收与系统内存回收
在worker函数末尾强制清空结果列表,触发Python垃圾回收,同时调用系统级内存回收(需对应权限):
def worker(args): # ... 原有处理逻辑 ... # 清理内存 result.clear() import gc gc.collect() # 触发Linux内存页回收(需root或对应权限) import os os.sync() os.popen('echo 3 > /proc/sys/vm/drop_caches') return [] # 返回空标识,不再传递大列表
3. 优化smaller_df结构,避免写时复制触发
若smaller_df在工作进程中被隐式修改(如类型转换),会触发写时复制导致内存飙升。可将其转换为不可变结构:
# 主进程预处理阶段转换 smaller_df = [frozenset(item) for item in smaller_df]
不可变结构不会触发写时复制,减少工作进程的额外内存占用。
4. 改用imap_unordered分批处理结果
pool.map会等待所有任务完成后一次性返回结果,改用imap_unordered可分批获取结果,主进程拿到一批就处理并释放内存:
with multiprocessing.Pool(processes=3, maxtasksperchild=5) as p: # chunksize控制每次分发的任务数,减少调度开销 for batch_result in p.imap_unordered(worker, chunk_indices, chunksize=10): # 立即处理结果,比如写入文件 process_and_write(batch_result) p.close() p.join()
内容的提问来源于stack exchange,提问作者Andrey

