【问题标题】:Writing a parallel loop编写并行循环
【发布时间】:2016-02-20 15:19:26
【问题描述】:

我正在尝试在一个简单的示例上运行并行循环。
我做错了什么?

from joblib import Parallel, delayed  
import multiprocessing

def processInput(i):  
        return i * i

if __name__ == '__main__':

    # what are your inputs, and what operation do you want to 
    # perform on each input. For example...
    inputs = range(1000000)      

    num_cores = multiprocessing.cpu_count()

    results = Parallel(n_jobs=4)(delayed(processInput)(i) for i in inputs) 

    print(results)

代码的问题在于,当在 Python 3 中的 Windows 环境下执行时,它会打开 num_cores 的 python 实例来执行并行作业,但只有一个处于活动状态。这不应该是这种情况,因为处理器的活动应该是 100% 而不是 14%(在 i7 - 8 逻辑内核下)。

为什么额外的实例没有做任何事情?

【问题讨论】:

  • 您是否收到任何错误消息?它对我来说运行良好......缩进应该是 4 个空格而不是 1 个...
  • 我也有同样的问题。问题是代码只在一个核上运行,而不是在 n 核上。

标签: python windows parallel-processing joblib


【解决方案1】:

继续根据您提供工作多处理代码的要求,我建议您使用pool_map(如果延迟功能不重要),我会给您举个例子,如果您正在使用 python3 值得提到你可以使用星图。 另外值得一提的是,如果返回结果的顺序不必与输入的顺序相对应,您可以使用 map_sync/starmap_async。

import multiprocessing as mp

def processInput(i):
        return i * i

if __name__ == '__main__':

    # what are your inputs, and what operation do you want to
    # perform on each input. For example...
    inputs = range(1000000)
    #  removing processes argument makes the code run on all available cores
    pool = mp.Pool(processes=4)
    results = pool.map(processInput, inputs)
    print(results)

【讨论】:

  • 我喜欢它的简单性,所以我试了一下。我得到一个 TypeError: cannot serialize '_io.TextIOWrapper' 对象。我的函数很复杂,我没有时间深入研究它,只是评论一下如果你有一个复杂的函数,这可能不是开箱即用的
  • 序列化是每个多进程程序的主要部分。为了尝试缓解此类问题,我建议检查您的复杂函数并检查它的哪一部分确实需要多处理解决方案,并尝试将其与复杂函数解耦,这将简化序列化,甚至可能使其变得不必要。
【解决方案2】:

在 Windows 上,多处理模块使用 'spawn' 方法来启动多个 python 解释器进程。这是相对缓慢的。 Parallel 试图聪明地运行代码。特别是,它会尝试调整批量大小,以便批量执行大约需要半秒。 (请参阅https://pythonhosted.org/joblib/parallel.html 处的 batch_size 参数)

您的 processInput() 函数运行速度如此之快,以至于 Parallel 确定在一个处理器上串行运行作业比启动多个 python 解释器并并行运行代码要快。

如果您想强制您的示例在多个内核上运行,请尝试将 batch_size 设置为 1000 或使 processInput() 更复杂,以便执行更长时间。

编辑:在 Windows 上显示多个正在使用的进程的工作示例(我使用的是 Windows 7):

from joblib import Parallel, delayed
from os import getpid

def modfib(n):
    # print the process id to see that multiple processes are used, and
    # re-used during the job.
    if n%400 == 0:
        print(getpid(), n)  

    # fibonacci sequence mod 1000000
    a,b = 0,1
    for i in range(n):
        a,b = b,(a+b)%1000000
    return b

if __name__ == "__main__":
    Parallel(n_jobs=-1, verbose=5)(delayed(modfib)(j) for j in range(1000, 4000))

【讨论】:

  • 您能否提议修改代码以使任务有效地并行执行?由于上面的代码是作为joblib使用的示例给出的,所以应该有一个实际有效的示例。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2016-03-08
  • 1970-01-01
  • 2015-06-04
  • 1970-01-01
  • 2020-03-19
  • 1970-01-01
  • 2020-06-13
相关资源
最近更新 更多