如何在Python starmap_async中实现多进程安全写文件记录OpenFOAM模拟状态
问题:OpenFOAM自动化模拟流水线的崩溃记录异常
我用Python的PyFoam库搭建OpenFOAM的自动化模拟流水线,用来生成机器学习所需的大型数据库(约50万次独立模拟)。为在多机器上运行流水线,采用multiprocessing.Pool.starmap_async(args)实现旧模拟完成后自动启动新模拟。但部分模拟可能崩溃,我需要生成文本文件记录所有崩溃案例。尝试使用multiprocessing.Manager.Queue()加监听器的方案后,在starmap_async()下无法正常运行——测试时所有模拟均成功,但文本文件仅写入一条完成案例记录。求问如何为每个完成的模拟(尤其是崩溃案例)向文件写入信息,相关代码片段如下:
from PyFoam.Execution.BasicRunner import BasicRunner from PyFoam.Execution.ParallelExecution import LAMMachine import numpy as np import multiprocessing import itertools import psutil # Defining global variables manager = multiprocessing.Manager() queue = manager.Queue() def runCase(airfoil, angle, velocity): # define simulation name newCase = str(airfoil) + "_" + str(angle) + "_" + str(velocity) ''' A lot of pre-processing commands to prepare the simulation which has been removed from snipped such as generate geometry, create mesh etc... ''' # run simulation machine = LAMMachine(nr=4) # set number of cores for parallel execution simulation = BasicRunner(argv=[solver, "-case", case.name], silent=True, lam=machine, logname="solver") simulation.start() # start simulation # check if simulation has completed if simulation.runOK(): # write message into queue queue.put(newCase) if not simulation.runOK(): print("Simulation did not run successfully") def listener(queue): fname = 'errors.txt' msg = queue.get() while True: with open(fname, 'w') as f: if msg == 'complete': break f.write(str(msg) + '\n') def main(): # Create parameter list angles = np.arange(-5, 0, 1) machs = np.array([0.15]) nacas = ['0012'] paramlist = list(itertools.product(nacas, angles, np.round(machs, 9))) # create number of processes and keep 2 cores idle for other processes nCores = psutil.cpu_count(logical=False) - 2 nProc = 4 nProcs = int(nCores / nProc) with multiprocessing.Pool(processes=nProcs) as pool: pool.apply_async(listener, (queue,)) # start the listener pool.starmap_async(runCase, paramlist).get() # run parallel simulations queue.put('complete') pool.close() pool.join() if __name__ == '__main__': main()
问题分析与修复方案
你的代码存在4个核心问题,导致监听器无法正常工作:
- 监听器仅读取一次消息:
listener函数仅在初始化时调用一次queue.get(),后续循环未获取新消息,只能处理第一条数据。 - 文件写入模式错误:用
'w'模式打开文件会清空原有内容,应该用'a'追加模式保留所有记录。 - 结束信号发送时机错误:
queue.put('complete')放在with multiprocessing.Pool块外,此时进程池已关闭,无法传递结束信号。 - 崩溃案例未记录:原代码仅将成功案例放入队列,未记录崩溃模拟,不符合需求。
修改后的完整代码
from PyFoam.Execution.BasicRunner import BasicRunner from PyFoam.Execution.ParallelExecution import LAMMachine import numpy as np import multiprocessing import itertools import psutil # 定义全局队列 manager = multiprocessing.Manager() queue = manager.Queue() def runCase(airfoil, angle, velocity): # 定义模拟名称 newCase = f"{airfoil}_{angle}_{velocity}" ''' 预处理命令已省略(如生成几何、划分网格等) ''' # 运行模拟 machine = LAMMachine(nr=4) # 设置并行核心数 simulation = BasicRunner(argv=[solver, "-case", case.name], silent=True, lam=machine, logname="solver") simulation.start() # 向队列发送结果:标记成功/失败+案例名 if simulation.runOK(): queue.put(f"SUCCESS: {newCase}") else: queue.put(f"FAILED: {newCase}") print(f"Simulation {newCase} did not run successfully") def listener(queue): fname = 'simulation_results.txt' # 用追加模式打开文件,避免清空已有内容 with open(fname, 'a') as f: while True: msg = queue.get() if msg == 'complete': break f.write(f"{msg}\n") # 立即刷新缓冲区,确保记录及时写入文件 f.flush() def main(): # 创建参数列表 angles = np.arange(-5, 0, 1) machs = np.array([0.15]) nacas = ['0012'] paramlist = list(itertools.product(nacas, angles, np.round(machs, 9))) # 计算可用进程数:保留2个物理核心,每个模拟用4核心 nCores = psutil.cpu_count(logical=False) - 2 nProcPerSim = 4 nProcs = int(nCores / nProcPerSim) with multiprocessing.Pool(processes=nProcs) as pool: # 启动监听器进程 pool.apply_async(listener, (queue,)) # 启动所有模拟并等待完成 pool.starmap_async(runCase, paramlist).get() # 发送结束信号给监听器 queue.put('complete') # with块结束后会自动调用pool.close()和pool.join() if __name__ == '__main__': main()
关键修改说明
- 循环读取队列消息:将
msg = queue.get()移至while循环内部,确保每次循环都获取新消息。 - 文件追加模式:用
'a'模式打开文件,每次写入都会追加到文件末尾,不会覆盖原有记录。 - 调整结束信号时机:将
queue.put('complete')放在with块内部,确保进程池运行时发送结束信号。 - 区分成功与失败记录:在队列消息中加入
SUCCESS/FAILED标记,方便后续筛选崩溃案例。 - 缓冲区实时刷新:调用
f.flush()确保写入内容立即保存到文件,避免因缓存导致记录丢失。
内容的提问来源于stack exchange,提问作者montju
相关产品推荐
相关产品推荐

