【问题标题】:What python pattern can be used to parallelization?什么python模式可以用于并行化?
【发布时间】:2018-11-20 12:23:45
【问题描述】:

cmd 是一个处理参数 x 的函数,将输出打印到标准输出。例如,它可能是

def cmd(x):
  print(x)

调用cmd() 的串行程序如下所示。

for x in array:
  cmd(x)

为了加快程序速度,我希望它并行运行。 stdout 输出可以是无序的,但单个 x 的输出不能被另一个 x 的输出破坏。

在 python 中可以有多种实现方式。我想出类似的东西。

from joblib import Parallel, delayed
Parallel(n_jobs=100)(delayed(cmd)(i) for i in range(100))

就代码的简单性/可读性和效率而言,这是在 python 中实现这一点的最佳方式吗?

另外,上面的代码在 python3 上运行正常。但不是在python2上,我收到以下错误。是不是可能导致错误的问题?

/Library/Frameworks/Python.framework/Versions/2.7/lib/python2.7/site-packages/joblib/externals/loky/backend/semlock.py:217:RuntimeWarning:信号量在 OSX 上被破坏,发布可能增加其最大值 “增加其最大值”,RuntimeWarning)

谢谢。

【问题讨论】:

  • 我修复了原始消息中的错误。

标签: python multithreading python-2.7 parallel-processing joblib


【解决方案1】:

在标准库https://docs.python.org/3/library/threading.html

import threading

def cmd(x):
    lock.acquire(blocking=True)
    print(x)
    lock.release()

lock = threading.Lock()

for i in range(100):
    t = threading.Thread(target=cmd, args=(i,))
    t.start()

使用锁可以保证lock.acquire()lock.release()之间的代码一次只被一个线程执行。 print 方法在 python3 中已经是线程安全的,因此即使没有锁也不会中断输出。但是,如果您在线程(它们修改的对象)之间共享任何状态,则需要一个锁。

【讨论】:

  • 如何确保同时运行的cmd()不超过n个?
  • 像我这样的简单方法是没有办法的。并行处理是一个复杂的课题。更高级的线程解决方案从这里开始:stackoverflow.com/questions/2846653/…
【解决方案2】:

如果您使用的是 python3,那么您可以使用标准库中的concurrent.futures

考虑以下用法:

with concurrent.futures.ProcessPoolExecutor(100) as executor:
     for x in array:
         executor.submit(cmd, x)

【讨论】:

  • 只是为了确定。这是否保证打印结果不会相互交错?
  • 如果您使用的是 python3,那么可以。但更好的解决方案可能是使用logging 模块
  • 结果会立即可用吗?当我在生产中使用代码时,我不确定为什么结果不会立即打印出来。我不确定这是由于 io 缓冲区问题还是与 ProcessPoolExecutor 相关的问题。冲洗 io Buffet 会导致问题吗?
  • 您的cmd 方法除了调用print 之外还有其他用途吗?
  • 是的。有一些代码进行计算,然后打印结果。
【解决方案3】:

我将使用以下代码解决问题中的问题(假设我们讨论的是 CPU 绑定操作):

import multiprocessing as mp
import random


def cmd(value):
    # some CPU heavy calculation
    for dummy in range(10 ** 8):
        random.random()
    # result
    return "result for {}".format(value)


if __name__ == '__main__':
    data = [val for val in range(10)]
    pool = mp.Pool(4)  # 4 - is the number of processes (the number of CPU cores used)
    # result is obtained after the process of all the data
    result = pool.map(cmd, data)

    print(result)

输出:

['result for 0', 'result for 1', 'result for 2', 'result for 3', 'result for 4', 'result for 5', 'result for 6', 'result for 7', 'result for 8', 'result for 9']

EDIT - 计算后立即获得结果的另一个实现 - processesqueues 而不是 poolmap

import multiprocessing
import random


def cmd(value, result_queue):
    # some CPU heavy calculation
    for dummy in range(10 ** 8):
        random.random()
    # result
    result_queue.put("result for {}".format(value))


if __name__ == '__main__':

    data = [val for val in range(10)]
    results = multiprocessing.Queue()

    LIMIT = 3  # 3 - is the number of processes (the number of CPU cores used)
    counter = 0
    for val in data:
        counter += 1
        multiprocessing.Process(
            target=cmd,
            kwargs={'value': val, 'result_queue': results}
            ).start()
        if counter >= LIMIT:
            print(results.get())
            counter -= 1
    for dummy in range(LIMIT - 1):
        print(results.get())

输出:

result for 0
result for 1
result for 2
result for 3
result for 4
result for 5
result for 7
result for 6
result for 8
result for 9

【讨论】:

  • 我想在结果可用时立即打印(无序也可以)。 pool.map()可以吗?
  • @user1424739 我想Pool.map() 是不可能的。检查编辑添加到另一个版本的代码的答案。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-02-23
  • 2016-10-22
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多