【问题标题】:Python: running subprocess in parallel [duplicate]Python:并行运行子进程[重复]
【发布时间】:2013-05-03 06:43:47
【问题描述】:

我有以下将 md5sum 写入日志文件的代码

for file in files_output:
    p=subprocess.Popen(['md5sum',file],stdout=logfile)
p.wait()
  1. 这些会并行编写吗?即如果 md5sum 需要很长时间来处理其中一个文件,是否会在等待前一个文件完成之前启动另一个文件?

  2. 1234563 (有些文件很大,有些很小)

【问题讨论】:

    标签: python subprocess


    【解决方案1】:

    从并行 md5sum 子进程收集输出的一种简单方法是使用线程池并从主进程写入文件:

    from multiprocessing.dummy import Pool # use threads
    from subprocess import check_output
    
    def md5sum(filename):
        try:
            return check_output(["md5sum", filename]), None
        except Exception as e:
            return None, e
    
    if __name__ == "__main__":
        p = Pool(number_of_processes) # specify number of concurrent processes
        with open("md5sums.txt", "wb") as logfile:
            for output, error in p.imap(md5sum, filenames): # provide filenames
                if error is None:
                   logfile.write(output)
    
    • md5sum 的输出很小,因此您可以将其存储在内存中
    • imap 保持秩序
    • number_of_processes 可能与文件或 CPU 内核的数量不同(值越大并不意味着更快:它取决于 IO(磁盘)和 CPU 的相对性能)

    您可以尝试一次将多个文件传递给 md5sum 子进程。

    在这种情况下,您不需要外部子流程; you can calculate md5 in Python:

    import hashlib
    from functools import partial
    
    def md5sum(filename, chunksize=2**15, bufsize=-1):
        m = hashlib.md5()
        with open(filename, 'rb', bufsize) as f:
            for chunk in iter(partial(f.read, chunksize), b''):
                m.update(chunk)
        return m.hexdigest()
    

    要使用多个进程而不是线程(以允许纯 Python md5sum() 使用多个 CPU 并行运行),只需从上述代码的导入中删除 .dummy

    【讨论】:

    • 抱歉,这里还在学习。我不明白为什么这里不使用队列。如果多个进程正在写入日志文件,会不会有问题?如果我弄错了,怎么同步?
    • 看起来Pool 支持异步调用。这是否意味着它保持 md5 的写入顺序(以filenames 的顺序)?不像简单地启动 x 个线程?
    • Pool 提供更高级别的接口。它在内部使用Queues 本身。 logfile 文件只能从主线程访问(只有md5sum() 函数在子线程中执行)。 imap() 按顺序返回结果(正如我已经明确提到的)
    • 你能指出一些可以帮助我学习论文主题的东西吗?我试过用谷歌搜索,但没有找到任何关于多处理的全面和介绍性的内容。
    • 这是一个很大的话题(尝试查找并发/并行/分布式编程/计算)。您对哪些特定方面感兴趣?
    【解决方案2】:
    1. 是的,这些 md5sum 进程将并行启动。
    2. 是的,md5sums 的写入顺序是不可预测的。通常,以这种方式从多个进程共享单个资源(如文件)被认为是一种不好的做法。

    您在for 循环之后创建p.wait() 的方式也将等待最后一个 md5sum 进程完成,而其余进程可能仍在运行。

    但是,如果您将 md5sum 输出收集到临时文件中并在所有处理完成后将其收集回一个文件中,您可以稍微修改此代码以仍然具有并行处理和同步输出的可预测性的好处。

    import subprocess
    import os
    
    processes = []
    for file in files_output:
        f = os.tmpfile()
        p = subprocess.Popen(['md5sum',file],stdout=f)
        processes.append((p, f))
    
    for p, f in processes:
        p.wait()
        f.seek(0)
        logfile.write(f.read())
        f.close()
    

    【讨论】:

    • 所以我猜这里的顺序是保留的,因为 processes[] 会跟踪它? IE。 process.append((p,f)) 在 md5sum 完成之前执行,按照 files_output 的顺序。
    • 是的,processes[] 将保留files_output[] 的原始顺序,并确保每个 md5sum 过程都完成。但是如果您担心操作系统的资源,您应该考虑使用任务队列和同步 md5sum 在每个线程中运行的线程池,@Alfe 建议使用subprocess.check_output()
    【解决方案3】:

    所有子进程并行运行。 (为了避免这种情况,必须明确等待它们的完成。)它们甚至可以同时写入日志文件,从而使输出混乱。为避免这种情况,您应该让每个进程写入不同的日志文件,并在所有进程完成后收集所有输出。

    q = Queue.Queue()
    result = {}  # used to store the results
    for fileName in fileNames:
      q.put(fileName)
    
    def worker():
      while True:
        fileName = q.get()
        if fileName is None:  # EOF?
          return
        subprocess_stuff_using(fileName)
        wait_for_finishing_subprocess()
        checksum = collect_md5_result_for(fileName)
        result[fileName] = checksum  # store it
    
    threads = [ threading.Thread(target=worker) for _i in range(20) ]
    for thread in threads:
      thread.start()
      q.put(None)  # one EOF marker for each thread
    

    在此之后,结果应存储在result

    【讨论】:

    • 谢谢。但是,我有 1000 个 md5sum。我宁愿不为每个文件打开一个单独的文件。
    • 不,你不应该。创建一个Queue.Queue和几十个线程的线程池,让每个线程从队列中读取一个元素并为这个元素启动一个子进程,等待这个子进程完成,得到结果(md5校验和),存储映射中的结果。如果队列为空,线程应该终止。
    • 仍然是 Python 新手。我是否需要使用 Queue.Queue 同时写入映射?如果没有,Queue.Queue 能为我做什么?
    • 不直接看你的代码我就不知道了。代码当前声明如下:将所有任务(按原始顺序)放入队列中,并告诉 20 个工作人员每个人这样做:从队列中取出一个任务并处理它,继续直到从队列中获得 EOF(无)。因为工人是并行工作的,这当然意味着最后一个完成他的第一个任务(第二十个任务)的工人可以最先完成他的任务。这将改变结果到达的顺序。但这取决于任务需要的时间。
    • 此代码中的结果存储在结果映射中,因此在此收集所有结果后,您可以遍历原始列表(按您想要的顺序)并得到匹配的结果结果字典。
    猜你喜欢
    • 2020-06-16
    • 2011-12-16
    • 1970-01-01
    • 2014-07-12
    • 1970-01-01
    • 2022-01-25
    • 2018-08-13
    • 1970-01-01
    • 2021-11-10
    相关资源
    最近更新 更多