【问题标题】:Using Python's Multiprocessing module to execute simultaneous and separate SEAWAT/MODFLOW model runs使用 Python 的 Multiprocessing 模块执行同时和单独的 SEAWAT/MODFLOW 模型运行
【发布时间】: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?

Python Multiprocessing queue

我的流程:

我想通过 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() 以从不同的目录开始。

标签: python multiprocessing


【解决方案1】:

我在 Python 代码中看不到任何计算。如果你只需要并行执行几个外部程序,使用subprocess 来运行程序和threading 模块来保持恒定数量的进程运行就足够了,但最简单的代码是使用multiprocessing.Pool:

#!/usr/bin/env python
import os
import multiprocessing as mp

def run(filename_def_param): 
    filename, def_param = filename_def_param # unpack arguments
    ... # call external program on `filename`

def safe_run(*args, **kwargs):
    """Call run(), catch exceptions."""
    try: run(*args, **kwargs)
    except Exception as e:
        print("error: %s run(*%r, **%r)" % (e, args, kwargs))

def main():
    # populate files
    ws = r'D:\Data\Users\jbellino\Project\stJohnsDeepening\model\xsec_a'
    workdir = os.path.join(ws, r'fieldgen\reals')
    files = ((os.path.join(workdir, f), ws)
             for f in os.listdir(workdir) if f.endswith('.npy'))

    # start processes
    pool = mp.Pool() # use all available CPUs
    pool.map(safe_run, files)

if __name__=="__main__":
    mp.freeze_support() # optional if the program is not frozen
    main()

如果文件很多,则pool.map() 可以替换为for _ in pool.imap_unordered(safe_run, files): pass。

还有mutiprocessing.dummy.Pool 提供与multiprocessing.Pool 相同的接口,但使用线程而不是在这种情况下可能更合适的进程。

您不需要保留一些 CPU 空闲。只需使用一个以低优先级启动可执行文件的命令(在 Linux 上它是一个nice 程序)。

ThreadPoolExecutor example

concurrent.futures.ThreadPoolExecutor 既简单又足够,但它需要3rd-party dependency on Python 2.x(它自 Python 3.2 起就在 stdlib 中)。

#!/usr/bin/env python
import os
import concurrent.futures

def run(filename, def_param):
    ... # call external program on `filename`

# populate files
ws = r'D:\Data\Users\jbellino\Project\stJohnsDeepening\model\xsec_a'
wdir = os.path.join(ws, r'fieldgen\reals')
files = (os.path.join(wdir, f) for f in os.listdir(wdir) if f.endswith('.npy'))

# start threads
with concurrent.futures.ThreadPoolExecutor(max_workers=8) as executor:
    future_to_file = dict((executor.submit(run, f, ws), f) for f in files)

    for future in concurrent.futures.as_completed(future_to_file):
        f = future_to_file[future]
        if future.exception() is not None:
           print('%r generated an exception: %s' % (f, future.exception()))
        # run() doesn't return anything so `future.result()` is always `None`

或者如果我们忽略run() 引发的异常:

from itertools import repeat

... # the same

# start threads
with concurrent.futures.ThreadPoolExecutor(max_workers=8) as executor:
     executor.map(run, files, repeat(ws))
     # run() doesn't return anything so `map()` results can be ignored

subprocess + threading(手动池)解决方案

#!/usr/bin/env python
from __future__ import print_function
import os
import subprocess
import sys
from Queue import Queue
from threading import Thread

def run(filename, def_param):
    ... # define exe, swt_nam
    subprocess.check_call([exe, swt_nam]) # run external program

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)

# start threads
q = Queue()
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()

# populate files
ws = r'D:\Data\Users\jbellino\Project\stJohnsDeepening\model\xsec_a'
wdir = os.path.join(ws, r'fieldgen\reals')
for f in os.listdir(wdir):
    if f.endswith('.npy'):
        q.put_nowait((os.path.join(wdir, f), ws))

for _ in threads: q.put_nowait(None) # signal no more files
for t in threads: t.join() # wait for completion

【讨论】:

  • 感谢您的回答,我宁愿坚持使用 MP 模块,因为我最近几天都在阅读它;如果我不需要的话,我现在不想切换到别的东西。然而,该函数同时调用所有 100 个实现 - 99 个立即关闭,我只剩下一个实际运行的。我想我曾经尝试过 Popen 模块并得到了类似的结果。有任何想法吗? mp.cpu_count = 8.
  • 再次运行它,看起来它正在调用模型以 7 个批次运行(我设置了 mp.Pool(processes=mp.cpu_count()-1)),但它不会等待这 7 个运行完成才调用下一批等等。进步!
  • @Jason: run() 函数必须阻塞,直到给定 filename 的所有工作完成。将os.system(exe + swt_nam) 替换为subprocess.check_call([exe, swt_nam])。它会产生任何错误吗?它是立即返回还是等待?检查所有路径是否正确。
  • 我得到了同样的行为,除了在我关闭剩下的运行的单独模型之后我得到这个错误:Exception in thread Thread-2: Traceback (most recent call last): File "C:\Python26\lib\threading.py", line 532, in __bootstrap_inner self.run() File "C:\Python26\lib\threading.py", line 484, in run self.__target(*self.__args, **self.__kwargs) File "C:\Python26\lib\multiprocessing\pool.py", line 259, in _handle_results task = get() TypeError: ('__init__() takes exactly 3 arguments (1 given)', <class 'subprocess.CalledProcessError'>, ())
  • PS - 我刚刚意识到脚本在抛出该错误后被挂起,不得不手动杀死它。
【解决方案2】:

这是我在内存中保持最小 x 线程数的方法。它是线程和多处理模块的组合。对于其他技术,如受人尊敬的成员在上面解释过,这可能是不寻常的,但可能非常值得。为了解释起见,我假设一次抓取至少 5 个网站。

所以这里是:-

#importing dependencies.
from multiprocessing import Process
from threading import Thread
import threading

# Crawler function
def crawler(domain):
    # define crawler technique here.
    output.write(scrapeddata + "\n")
    pass

接下来是threadController函数。该函数将控制线程流向主内存。它将继续激活线程以维持 threadNum“最小”限制,即。 5. 直到所有活动线程(acitveCount)都完成后才会退出。

它将保持最少的 threadNum(5) startProcess 函数线程(这些线程最终会从 processList 启动进程,同时在 60 秒内加入它们)。启动 threadController 后,将有 2 个线程不包括在上述 5 个限制中,即。 Main 线程和 threadController 线程本身。这就是为什么使用 threading.activeCount() != 2 的原因。

def threadController():
    print "Thread count before child thread starts is:-", threading.activeCount(), len(processList)
    # staring first thread. This will make the activeCount=3
    Thread(target = startProcess).start()
    # loop while thread List is not empty OR active threads have not finished up.
    while len(processList) != 0 or threading.activeCount() != 2:
        if (threading.activeCount() < (threadNum + 2) and # if count of active threads are less than the Minimum AND
            len(processList) != 0):                            # processList is not empty
                Thread(target = startProcess).start()         # This line would start startThreads function as a seperate thread **

startProcess 函数作为一个单独的线程,将从进程列表中启动进程。这个函数的目的(**作为一个不同的线程开始)是它将成为进程的父线程。因此,当它将以 60 秒的超时时间加入它们时,这将停止 startProcess 线程继续前进,但这不会停止 threadController 执行。这样一来,threadController 就会按要求工作。

def startProcess():
    pr = processList.pop(0)
    pr.start()
    pr.join(60.00) # joining the thread with time out of 60 seconds as a float.

if __name__ == '__main__':
    # a file holding a list of domains
    domains = open("Domains.txt", "r").read().split("\n")
    output = open("test.txt", "a")
    processList = [] # thread list
    threadNum = 5 # number of thread initiated processes to be run at one time

    # making process List
    for r in range(0, len(domains), 1):
        domain = domains[r].strip()
        p = Process(target = crawler, args = (domain,))
        processList.append(p) # making a list of performer threads.

    # starting the threadController as a seperate thread.
    mt = Thread(target = threadController)
    mt.start()
    mt.join() # won't let go next until threadController thread finishes.

    output.close()
    print "Done"

除了在内存中保持最少数量的线程外,我的目标是还有一些东西可以避免内存中的线程或进程卡住。我使用超时功能做到了这一点。 对于任何打字错误,我深表歉意。

我希望这个结构可以帮助这个世界上的任何人。 问候, 维卡斯·高塔姆

【讨论】:

    猜你喜欢
    • 2016-12-04
    • 1970-01-01
    • 2011-12-20
    • 1970-01-01
    • 2015-09-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多