【问题标题】:Correct way of using multiprocessing Process() for parallel execution使用 multiprocessing Process() 进行并行执行的正确方法
【发布时间】:2017-05-31 15:21:39
【问题描述】:

如果我的概念在使用 Python 的多处理并行执行 exe 文件时是否存在根本性错误,我能否与你们核实一下。

所以我有一大堆工作(示例代码中为 100000 个),我想使用所有可用的内核(我的计算机中为 16 个)并行运行它们。下面的代码没有像我看到的许多示例那样使用队列,但它似乎有效。只是想避免代码“工作”的情况,但是当我将其扩展到运行多个计算节点时,等待炸毁存在一个巨大的错误。有人可以帮忙吗?

import subprocess
import multiprocessing

def task_fn(task_dir) :
    cmd_str = ["my_exe","-my_exe_arguments"]
    try :
        msg = subprocess.check_output(cmd_str,cwd=task_dir,stderr=subprocess.STDOUT,universal_newlines=True)
    except subprocess.CalledProcessError as e :
        with open("a_unique_err_log_file.log","w") as f :
            f.write(e.output)
    return;

if __name__ == "__main__":

    n_cpu = multiprocessing.cpu_count()
    num_jobs = 100000
    proc_list = [multiprocessing.Process() for p in range(n_cpu)]

    for i in range(num_jobs):
        task_dir = str(i)
        task_processed = False
        while not(task_processed) :
            # Search through all processes in p_list repeatedly until a 
            # terminated processs is found to take on a new task
            for p in range(len(p_list)) :
                if not(p_list[p].is_alive()) :
                    p_list[p] = multiprocessing.Process(target=task_fn,args=(task_dir,))
                    p_list[p].start()
                    task_processed = True

    # At the end of the outermost for loop 
    # Wait until all the processes have finished
    for p in p_list :
        p.join()

    print("All Done!")

【问题讨论】:

    标签: python-3.x parallel-processing subprocess multiprocessing


    【解决方案1】:

    与其自己生成和管理进程,不如使用Pool of workers。它旨在为您处理所有这些问题。

    当您的工作人员正在生成子进程时,您可以使用线程而不是进程。

    此外,工人似乎会写在同一个文件上。您需要保护其访问不受并发实例的影响,否则结果将完全无序。

    from threading import Lock
    from concurrent.futures import ThreadPoolExecutor  
    
    
    mutex = Lock()
    task_dir = "/tmp/tasks"
    
    
    def task_fn(task_nr):  
        """This function will run in a separate thread."""
        cmd_str = ["my_exe","-my_exe_arguments"]
        try:
            msg = subprocess.check_output(cmd_str, cwd=task_dir, stderr=subprocess.STDOUT, universal_newlines=True)
        except subprocess.CalledProcessError as e:
            with mutex:
                with open("a_unique_PROTECTED_err_log_file.log", "w") as f :
                    f.write(e.output)
    
        return task_nr
    
    
    with ThreadPoolExecutor() as pool:
        iterator = pool.map(task_fn, range(100000))
        for result in iterator:
            print("Task %d done" % result)
    

    【讨论】:

    • 您好,感谢您的回复!会试试的。但是所以基本上我的编码方式没有任何问题?
    • 是的。事实上,如果出现错误,您的日志文件将由于不受保护的文件访问而被全部打乱。而且,使用进程存在矫枉过正。只需使用线程。
    • 好吧,也许有一个误解,“a_unique_err_log_file.log”是我懒惰的说法,每个日志文件都是唯一的。这并不是字面意思,因为文件名实际上是“a_unique_err_log_file.log”,因此每个进程实际上都会写入一个唯一的文件。
    • 进程过大,对系统有什么影响?需要更多 CPU/内存资源?
    • 在 Unix 上,在您的具体情况下,影响应该可以忽略不计。在 Windows 上,每个进程都会有一个相当多的专用内存地址空间。在 Windows 上生成进程也更重。线程在两个平台上的内存和创建时间都是轻量级的。
    猜你喜欢
    • 2017-01-12
    • 2022-10-25
    • 2013-10-08
    • 2014-11-09
    • 2020-02-27
    • 1970-01-01
    • 2012-07-04
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多