如何基于空闲内存与函数参数动态控制Python多进程数量?
问题排查与解决方案
你的代码运行极快且未执行compute函数,核心原因是程序结束时没有等待所有提交的非大内存任务完成,ProcessPoolExecutor在主进程退出时会立即终止所有子进程,导致那些任务还没来得及执行。另外还有几个细节问题需要调整,下面一步步解决:
核心问题分析
- 未等待非大内存任务完成:你把非大内存任务的
future加入futures列表,但程序走到计时结束步骤时直接退出,没有等待这些future执行完毕,进程池被强制关闭,任务根本没机会运行。 futures列表未清空:处理大内存任务时,concurrent.futures.wait(futures)等待完现有任务后,没有清空列表,后续遇到大内存任务会重复等待已完成的任务,虽然不影响执行,但会造成冗余等待。- 潜在的
large_memory函数问题:如果这个函数的判断逻辑有误(比如始终返回False),所有任务都被加入futures,但最终没被等待,直接退出,这也会导致compute没被调用。
修改后的完整代码
import itertools import concurrent.futures import time import os # 请确保你的My_class、large_memory、compute函数已正确定义 # class My_class: # def __init__(self, x): # self.x = x # def large_memory(Obj, s_step, max_s): # # 替换为你实际的大内存任务判断逻辑 # return any(obj.x >= 50 for obj in Obj) # def compute(Obj, t0, tf, s_step, max_s, folder): # # 替换为你的实际计算逻辑 # print(f"正在计算: {[obj.x for obj in Obj]}") # Parameters N = int(input("Number of CPUs to use: ")) t0 = 0 tf = 200 s_step = 0.05 max_s = None folder = "test" possible_dynamics = [My_class(x) for x in [20, 30, 40, 50, 60]] dynamics_to_compute = [list(x) for x in itertools.combinations_with_replacement(possible_dynamics , 2)] + [list(x) for x in itertools.combinations_with_replacement(possible_dynamics , 3)] function_inputs = [(dyn , t0, tf, s_step, max_s, folder) for dyn in dynamics_to_compute] # ----------- # Computation # ----------- start = time.time() futures = [] # 使用上下文管理器管理进程池,确保自动清理资源 with concurrent.futures.ProcessPoolExecutor(max_workers=N) as pool: for Obj, t0, tf, s_step, max_s, folder in function_inputs: if large_memory(Obj, s_step, max_s): # 等待所有正在运行/待运行的非大内存任务完成 concurrent.futures.wait(futures) # 清空futures列表,避免后续重复等待已完成的任务 futures.clear() # 提交大内存任务并等待完成(单进程执行) large_future = pool.submit(compute, Obj, t0, tf, s_step, max_s, folder) large_future.result() else: # 提交非大内存任务并保存future对象 future = pool.submit(compute, Obj, t0, tf, s_step, max_s, folder) futures.append(future) # 等待所有剩余的非大内存任务完成 concurrent.futures.wait(futures) end = time.time() if round(end-start, 3) < 60: print ("Complete - Elapsed time: {} s".format(round(end-start,3))) else: print ("Complete - Elapsed time: {} mn and {} s".format(int((end-start)//60), round((end-start)%60,3))) os.system("pause")
关键修改点说明
- 补全必要导入:添加了
time和os模块,确保计时和终端暂停命令正常工作。 - 用上下文管理器管理进程池:
with concurrent.futures.ProcessPoolExecutor(...) as pool会在代码块结束时自动等待所有任务完成并安全关闭进程池,避免资源泄漏。 - 等待剩余非大内存任务:在循环结束后添加
concurrent.futures.wait(futures),确保所有非大内存任务都执行完毕再结束程序。 - 清空
futures列表:处理完大内存任务后清空列表,避免后续重复等待已完成的任务。 - 校验
large_memory逻辑:你需要保证这个函数能准确识别高内存消耗的任务组合,这是实现单进程执行大任务的核心前提。
额外优化建议
- 动态内存监控:如果需要根据实时空闲内存调整任务提交,可以使用
psutil库(先安装:pip install psutil),添加内存检查逻辑:
然后在提交非大内存任务前,等待内存充足:import psutil def enough_free_memory(threshold_gb=10): free_mem_gb = psutil.virtual_memory().available / (1024**3) return free_mem_gb >= threshold_gbelse: # 每10秒检查一次,直到有足够空闲内存再提交任务 while not enough_free_memory(): time.sleep(10) future = pool.submit(compute, Obj, t0, tf, s_step, max_s, folder) futures.append(future) - 任务排序优化:如果大内存任务较多,可以将它们单独提取出来放在最后执行,减少频繁等待任务完成的开销。
内容的提问来源于stack exchange,提问作者Mathieu
相关产品推荐
相关产品推荐

