【发布时间】:2012-04-10 01:40:22
【问题描述】:
我正在尝试在我的 8 处理器 64 位 Windows 7 机器上运行 100 个模型。我想同时运行 7 个模型实例以减少我的总运行时间(每个模型运行大约 9.5 分钟)。我已经查看了与 Python 的多处理模块有关的几个线程,但仍然缺少一些东西。
Using the multiprocessing module
How to spawn parallel child processes on a multi-processor system?
我的流程:
我想通过 SEAWAT/MODFLOW 运行 100 个不同的参数集以比较结果。我已经为每个模型运行预先构建了模型输入文件并将它们存储在它们自己的目录中。我想做的是一次运行 7 个模型,直到所有实现都完成。进程之间不需要通信或结果显示。到目前为止,我只能按顺序生成模型:
import os,subprocess
import multiprocessing as mp
ws = r'D:\Data\Users\jbellino\Project\stJohnsDeepening\model\xsec_a'
files = []
for f in os.listdir(ws + r'\fieldgen\reals'):
if f.endswith('.npy'):
files.append(f)
## def work(cmd):
## return subprocess.call(cmd, shell=False)
def run(f,def_param=ws):
real = f.split('_')[2].split('.')[0]
print 'Realization %s' % real
mf2k = r'c:\modflow\mf2k.1_19\bin\mf2k.exe '
mf2k5 = r'c:\modflow\MF2005_1_8\bin\mf2005.exe '
seawatV4 = r'c:\modflow\swt_v4_00_04\exe\swt_v4.exe '
seawatV4x64 = r'c:\modflow\swt_v4_00_04\exe\swt_v4x64.exe '
exe = seawatV4x64
swt_nam = ws + r'\reals\real%s\ss\ss.nam_swt' % real
os.system( exe + swt_nam )
if __name__ == '__main__':
p = mp.Pool(processes=mp.cpu_count()-1) #-leave 1 processor available for system and other processes
tasks = range(len(files))
results = []
for f in files:
r = p.map_async(run(f), tasks, callback=results.append)
我将if __name__ == 'main': 更改为以下内容,希望它能解决我认为for loop 在上述脚本中传递的并行性不足的问题。但是,模型甚至无法运行(没有 Python 错误):
if __name__ == '__main__':
p = mp.Pool(processes=mp.cpu_count()-1) #-leave 1 processor available for system and other processes
p.map_async(run,((files[f],) for f in range(len(files))))
非常感谢任何和所有帮助!
编辑 2012 年 3 月 26 日 13:31 EST
在@J.F. 中使用“手动池”方法。 Sebastian 在下面的回答我得到了我的外部 .exe 的并行执行。模型实现一次调用 8 个批次,但它不会等待这 8 个运行完成,然后再调用下一个批次,依此类推:
from __future__ import print_function
import os,subprocess,sys
import multiprocessing as mp
from Queue import Queue
from threading import Thread
def run(f,ws):
real = f.split('_')[-1].split('.')[0]
print('Realization %s' % real)
seawatV4x64 = r'c:\modflow\swt_v4_00_04\exe\swt_v4x64.exe '
swt_nam = ws + r'\reals\real%s\ss\ss.nam_swt' % real
subprocess.check_call([seawatV4x64, swt_nam])
def worker(queue):
"""Process files from the queue."""
for args in iter(queue.get, None):
try:
run(*args)
except Exception as e: # catch exceptions to avoid exiting the
# thread prematurely
print('%r failed: %s' % (args, e,), file=sys.stderr)
def main():
# populate files
ws = r'D:\Data\Users\jbellino\Project\stJohnsDeepening\model\xsec_a'
wdir = os.path.join(ws, r'fieldgen\reals')
q = Queue()
for f in os.listdir(wdir):
if f.endswith('.npy'):
q.put_nowait((os.path.join(wdir, f), ws))
# start threads
threads = [Thread(target=worker, args=(q,)) for _ in range(8)]
for t in threads:
t.daemon = True # threads die if the program dies
t.start()
for _ in threads: q.put_nowait(None) # signal no more files
for t in threads: t.join() # wait for completion
if __name__ == '__main__':
mp.freeze_support() # optional if the program is not frozen
main()
没有可用的错误回溯。 run() 函数在调用单个模型实现文件时执行其职责,就像调用多个文件一样。唯一的区别是,对于多个文件,它会被调用len(files) 次,尽管每个实例都会立即关闭,并且只允许完成一个模型运行,此时脚本会正常退出(退出代码 0)。
在main() 中添加一些打印语句会显示一些关于活动线程数和线程状态的信息(请注意,这只是对 8 个实现文件的测试,以使屏幕截图更易于管理,理论上所有 8 个文件应该可以同时运行,但是行为会在它们产生的地方继续并立即死亡,除了一个):
def main():
# populate files
ws = r'D:\Data\Users\jbellino\Project\stJohnsDeepening\model\xsec_a'
wdir = os.path.join(ws, r'fieldgen\test')
q = Queue()
for f in os.listdir(wdir):
if f.endswith('.npy'):
q.put_nowait((os.path.join(wdir, f), ws))
# start threads
threads = [Thread(target=worker, args=(q,)) for _ in range(mp.cpu_count())]
for t in threads:
t.daemon = True # threads die if the program dies
t.start()
print('Active Count a',threading.activeCount())
for _ in threads:
print(_)
q.put_nowait(None) # signal no more files
for t in threads:
print(t)
t.join() # wait for completion
print('Active Count b',threading.activeCount())
**“D:\\Data\\Users...”这行是我手动停止模型运行到完成时抛出的错误信息。一旦我停止模型运行,就会报告剩余的线程状态行并退出脚本。
编辑 2012 年 3 月 26 日 16:24 EST
SEAWAT 确实允许并发执行,就像我过去所做的那样,使用 iPython 手动生成实例并从每个模型文件夹启动。这一次,我从一个位置启动所有模型运行,即我的脚本所在的目录。看起来罪魁祸首可能在于 SEAWAT 保存部分输出的方式。当 SEAWAT 运行时,它会立即创建与模型运行相关的文件。这些文件之一没有保存到模型实现所在的目录中,而是保存在脚本所在的顶级目录中。这可以防止任何后续线程将相同的文件名保存在相同的位置(它们都希望这样做,因为这些文件名是通用的并且对每个实现都不是特定的)。 SEAWAT 窗口打开的时间不够长,我无法阅读甚至看到有错误消息,我只是在返回并尝试使用 iPython 运行代码时才意识到这一点,该代码直接显示来自 SEAWAT 的打印输出,而不是打开运行程序的新窗口。
我接受@J.F. Sebastian 的回答,因为很可能一旦我解决了这个模型可执行问题,他提供的线程代码将把我带到我需要的地方。
最终代码
在 subprocess.check_call 中添加 cwd 参数以在其自己的目录中启动每个 SEAWAT 实例。非常关键。
from __future__ import print_function
import os,subprocess,sys
import multiprocessing as mp
from Queue import Queue
from threading import Thread
import threading
def run(f,ws):
real = f.split('_')[-1].split('.')[0]
print('Realization %s' % real)
seawatV4x64 = r'c:\modflow\swt_v4_00_04\exe\swt_v4x64.exe '
cwd = ws + r'\reals\real%s\ss' % real
swt_nam = ws + r'\reals\real%s\ss\ss.nam_swt' % real
subprocess.check_call([seawatV4x64, swt_nam],cwd=cwd)
def worker(queue):
"""Process files from the queue."""
for args in iter(queue.get, None):
try:
run(*args)
except Exception as e: # catch exceptions to avoid exiting the
# thread prematurely
print('%r failed: %s' % (args, e,), file=sys.stderr)
def main():
# populate files
ws = r'D:\Data\Users\jbellino\Project\stJohnsDeepening\model\xsec_a'
wdir = os.path.join(ws, r'fieldgen\reals')
q = Queue()
for f in os.listdir(wdir):
if f.endswith('.npy'):
q.put_nowait((os.path.join(wdir, f), ws))
# start threads
threads = [Thread(target=worker, args=(q,)) for _ in range(mp.cpu_count()-1)]
for t in threads:
t.daemon = True # threads die if the program dies
t.start()
for _ in threads: q.put_nowait(None) # signal no more files
for t in threads: t.join() # wait for completion
if __name__ == '__main__':
mp.freeze_support() # optional if the program is not frozen
main()
【问题讨论】:
-
鉴于您的
run函数实际上会产生一个进程来完成这项工作,您也可以使用多线程而不是多处理。 -
感谢您的建议,如果我无法继续使用 MP 模块,我可能会走那条路 - 我不愿意切换到不同的模块,因为我已经沉没了这么多时间阅读这篇文章。
-
目前行为与预期行为有何不同尚不清楚。什么是预期行为?如果将
seawatV4x64调用替换为print_args.py会发生什么?顺便说一句,您不需要在threading解决方案中导入multiprocessing。 -
@J.F.Sebastian,预期的行为是代码为它在目录
fieldgen\reals中找到的每个参数文件运行一次模型。它将与mp.cpu_count()在自己的处理器上同时运行的模型数量并行执行此操作,直到所有参数文件都已运行。现在发生的情况是代码同时为所有参数文件生成所有模型运行,其中除了一个立即退出,我只剩下一个完整的模型运行。 -
您可以将
cwd=unique_for_the_model_directory参数添加到check_call()以从不同的目录开始。