【问题标题】:Multiprocessing thousands of files with external command使用外部命令多处理数千个文件
【发布时间】:2017-03-16 12:29:31
【问题描述】:

我想从 Python 为大约 8000 个文件启动一个外部命令。每个文件都独立于其他文件进行处理。唯一的限制是在处理完所有文件后继续执行。我有 4 个物理核心,每个都有 2 个逻辑核心(multiprocessing.cpu_count() 返回 8)。我的想法是使用一个由四个并行独立进程组成的池,这些进程将在 8 个内核中的 4 个上运行。这样我的机器应该可以同时使用。

这是我一直在做的事情:

import multiprocessing
import subprocess
import os
from multiprocessing.pool import ThreadPool


def process_files(input_dir, output_dir, option):
    pool = ThreadPool(multiprocessing.cpu_count()/2)
    for filename in os.listdir(input_dir):  # about 8000 files
        f_in = os.path.join(input_dir, filename)
        f_out = os.path.join(output_dir, filename)
        cmd = ['molconvert', option, f_in, '-o', f_out]
        pool.apply_async(subprocess.Popen, (cmd,))
    pool.close()
    pool.join()


def main():
    process_files('dir1', 'dir2', 'mol:H')
    do_some_stuff('dir2')
    process_files('dir2', 'dir3', 'mol:a')
    do_more_stuff('dir3')

对一批 100 个文件进行顺序处理需要 120 秒。上面概述的多处理版本(函数process_files)批处理只需要 20 秒。但是,当我对整组 8000 个文件运行 process_files 时,我的电脑会挂起并且一小时后没有解冻。

我的问题是:

1) 我认为ThreadPool 应该初始化一个进程池(确切地说,这里是multiprocessing.cpu_count()/2 进程)。但是,我的计算机挂断了 8000 个文件而不是 100 个文件,这表明可能没有考虑到池的大小。要么,要么我做错了什么。你能解释一下吗?

2) 当每个独立进程都必须启动外部命令时,这是在 Python 下启动独立进程的正确方法,并且所有资源都不会被处理占用吗?

【问题讨论】:

  • 我比较了@larsks(ThreadPoolapply_async 和子进程calls)和@Roland Smith(使用Popen 对象的手动池管理)提出的解决方案。我的基准测试表明ThreadPool 解决方案在实践中更快。非常感谢你们!

标签: python subprocess multiprocessing threadpool


【解决方案1】:

我认为您的基本问题是subprocess.Popen 的使用。该方法在返回之前等待命令完成。由于函数会立即返回(即使命令仍在运行),所以就您的 ThreadPool 而言,该函数已完成,并且它可以产生另一个......这意味着您最终会产生 8000 个左右的进程。

使用subprocess.check_call 可能会更好:

Run command with arguments.  Wait for command to complete.  If
the exit code was zero then return, otherwise raise
CalledProcessError.  The CalledProcessError object will have the
return code in the returncode attribute.

所以:

def process_files(input_dir, output_dir, option):
    pool = ThreadPool(multiprocessing.cpu_count()/2)
    for filename in os.listdir(input_dir):  # about 8000 files
        f_in = os.path.join(input_dir, filename)
        f_out = os.path.join(output_dir, filename)
        cmd = ['molconvert', option, f_in, '-o', f_out]
        pool.apply_async(subprocess.check_call, (cmd,))
    pool.close()
    pool.join()

如果您真的不关心退出代码,那么您可能需要subprocess.call,它不会在进程退出代码非零的情况下引发异常。

【讨论】:

  • 感谢您提供如此清晰和直截了当的解释。确实subprocess.Popen 一定是导致产生这么多进程的原因。我没有使用subprocess.call 认为 Python 会等待该过程完成,而不是用有用的工作人员填充池。但这就是为什么游泳池首先存在的原因。 (对不起,代表太低,无法投票。)
  • 您仍然可以通过单击此答案左侧的复选标记将其标记为“已接受”答案。
  • 是的,我知道。问题是我很难在两个非常有用的答案之间做出决定(我目前正在根据两个建议的解决方案测试结果)。 :D
【解决方案2】:

如果您使用的是 Python 3,我会考虑使用 concurrent.futures.ThreadPoolExecutormap 方法。

或者,您可以自己管理子流程列表。

下面的例子定义了一个函数来启动ffmpeg 来将视频文件转换为 Theora/Vorbis 格式。它为每个启动的子进程返回一个 Popen 对象。

def startencoder(iname, oname, offs=None):
    args = ['ffmpeg']
    if offs is not None and offs > 0:
        args += ['-ss', str(offs)]
    args += ['-i', iname, '-c:v', 'libtheora', '-q:v', '6', '-c:a',
            'libvorbis', '-q:a', '3', '-sn', oname]
    with open(os.devnull, 'w') as bb:
        p = subprocess.Popen(args, stdout=bb, stderr=bb)
    return p

在主程序中,代表正在运行的子进程的Popen对象列表是这样维护的。

outbase = tempname()
ogvlist = []
procs = []
maxprocs = cpu_count()
for n, ifile in enumerate(argv):
    # Wait while the list of processes is full.
    while len(procs) == maxprocs:
        manageprocs(procs)
    # Add a new process
    ogvname = outbase + '-{:03d}.ogv'.format(n + 1)
    procs.append(startencoder(ifile, ogvname, offset))
    ogvlist.append(ogvname)
# All jobs have been submitted, wail for them to finish.
while len(procs) > 0:
    manageprocs(procs)

因此,只有在运行的子进程少于核心时才会启动新进程。多次使用的代码被分离到manageprocs函数中。

def manageprocs(proclist):
    for pr in proclist:
        if pr.poll() is not None:
            proclist.remove(pr)
    sleep(0.5)

sleep 的调用用于防止程序在循环中旋转。

【讨论】:

  • 感谢您提及concurrent.futures.ThreadPoolExecutor(此处仍使用 Python 2.7)。感谢您提供这个手动池管理的好例子。我曾尝试做类似的事情(认为在迭代列表时我不应该在列表上做remove),但一定出了问题。我将很快测试这个解决方案。 (抱歉,rep 太低,无法投票。)
  • 我比较了这两种方法(你的答案和@larsks')。我非常喜欢这个解决方案,但似乎手动管理池会产生开销,这可能是由于对sleep 的调用(我让进程管理器休眠了 0.2 秒,因为它似乎更合适)。在我实际输入大小的 1/10 的批量测试中,手动池管理在 cpu_count()-1 内核上比 ThreadPool 慢 8%,在 cpu_count()/2 内核上比 ThreadPool 慢 27%。
  • 你必须做一些真正的分析才能看到差异来自哪里。影响事物的因素很多。例如,cpu_count() 不是给定的最佳子进程数量。您可能应该尝试从cpu_count()/2cpu_count()*2 范围内的任何东西。此外,您可能应该根据molconvert 通常花费的时间来调整sleep 的数量。但由于我已经完全切换到 Python 3,所以这些天我只是倾向于使用 concurrent.futures.ThreadPoolExecutor 来处理这样的事情。
猜你喜欢
  • 2017-11-30
  • 1970-01-01
  • 2015-08-03
  • 2012-08-09
  • 2015-05-03
  • 2015-07-29
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多