使用concurrent.futures.ThreadPoolExecutor主进程内存持续上涨不释放求助
我现在用concurrent.futures.ThreadPoolExecutor来并发执行代码,但运行一段时间后主进程内存占用一直涨,还释放不了。我的使用方式大概是这样:
with concurrent.futures.ThreadPoolExecutor(max_workers=10) as pool: processing_data = pool.map(handle_item, items_dict.values()) # 使用processing_data中的数据...
我查过相关问题,但场景都和我不太一样,那些解决方案似乎不适用。下面是能复现潜在内存泄漏的示例代码:
#!/usr/bin/env python3 import os import time import psutil import random import concurrent.futures from memory_profiler import profile as mem_profile p = psutil.Process(os.getpid()) def do_magic(values): return None @mem_profile def foo(): a = {i: chr(i) for i in range(1024)} with concurrent.futures.ThreadPoolExecutor(max_workers=10) as pool: proccessed_data = pool.map(do_magic, [1,2,3,4,5,6,7,8,9,10]) def fooer(): while True: foo() time.sleep(1) fooer()
有没有人知道我哪里错了,或者有什么解决办法?谢谢各位!
问题根源分析
你遇到的这个问题其实和ThreadPoolExecutor的线程复用机制有关。当你用with语句创建线程池时,它会在内部维护一组线程——这些线程不会在每次with块结束后就销毁,而是被保留下来复用(这本来是为了避免频繁创建销毁线程的性能开销)。
但麻烦的是,这些线程会持有之前任务相关对象的引用,哪怕任务已经执行完毕。尤其是当你的任务函数涉及大对象时,这些残留的引用会阻止Python的垃圾回收(GC)机制及时回收内存,久而久之就出现了内存持续上涨的情况。
看你的示例代码:每次调用foo()都会创建字典a,虽然看起来它和线程任务无关,但线程池里的线程在执行任务时,内部栈帧或局部变量可能会间接持有当前上下文的引用,再加上Python GC处理线程持有的引用时效率偏低,就导致了内存无法及时释放的假象。
可行的解决方法
这里给你几个实用方案,你可以根据自己的业务场景选择:
显式触发垃圾回收
在用完处理结果后,手动删除引用并强制触发GC,帮它“加速”回收内存。修改你的代码如下:import gc with concurrent.futures.ThreadPoolExecutor(max_workers=10) as pool: processing_data = pool.map(handle_item, items_dict.values()) # 使用processing_data中的数据... del processing_data # 显式删除结果引用 gc.collect() # 强制触发垃圾回收这个方法简单直接,适合快速验证问题是否由GC回收不及时导致。
优化任务函数,清理无效引用
检查你的任务函数(比如handle_item),确保它不会持有大对象的冗余引用,或者在任务结束后主动清理内部的大变量:def handle_item(item): # 处理逻辑 big_obj = some_heavy_operation(item) result = process_result(big_obj) del big_obj # 显式删除大对象引用 return result减少线程持有的无效引用,能让GC更高效地回收内存。
按需禁用线程复用(应急方案)
如果你不需要线程复用的性能优势,可以手动创建并关闭线程池,确保任务完成后销毁所有线程:pool = concurrent.futures.ThreadPoolExecutor(max_workers=10) processing_data = pool.map(handle_item, items_dict.values()) pool.shutdown(wait=True) # 等待所有任务完成后销毁线程 del pool # 删除线程池引用 # 使用processing_data中的数据...注意:这种方式会失去线程复用的性能增益,频繁创建销毁线程会增加开销,只适合任务执行频率不高的场景。
换成ProcessPoolExecutor(CPU密集型任务首选)
如果你的任务是CPU密集型,且不需要线程间共享内存,换成ProcessPoolExecutor可能是最优解。因为每个子进程有独立的内存空间,任务结束后系统会直接回收子进程的内存,不会出现主进程内存泄漏的问题:with concurrent.futures.ProcessPoolExecutor(max_workers=10) as pool: processing_data = pool.map(handle_item, items_dict.values())不过要注意,进程间通信有额外开销,大对象需要序列化传递,所以要根据任务类型判断是否适用。
针对你示例代码的验证方案
你可以在foo()函数结束前加上引用清理和GC触发,看看内存是否还会持续上涨:
@mem_profile def foo(): a = {i: chr(i) for i in range(1024)} with concurrent.futures.ThreadPoolExecutor(max_workers=10) as pool: proccessed_data = pool.map(do_magic, [1,2,3,4,5,6,7,8,9,10]) del a del proccessed_data import gc gc.collect()
修改后内存应该就能稳定下来了。
内容的提问来源于stack exchange,提问作者Or Y

