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

如何在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个核心问题,导致监听器无法正常工作:

  1. 监听器仅读取一次消息:listener函数仅在初始化时调用一次queue.get(),后续循环未获取新消息,只能处理第一条数据。
  2. 文件写入模式错误:用'w'模式打开文件会清空原有内容,应该用'a'追加模式保留所有记录。
  3. 结束信号发送时机错误:queue.put('complete')放在with multiprocessing.Pool块外,此时进程池已关闭,无法传递结束信号。
  4. 崩溃案例未记录:原代码仅将成功案例放入队列,未记录崩溃模拟,不符合需求。

修改后的完整代码

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 01:01:05