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

如何强制释放Python进程内存以避免累积分配引发MemoryError?

内存泄漏问题排查与解决方案(Python多进程处理大CSV场景)

环境信息

  • 操作系统: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 04:04:51