【问题标题】:How to parallelise code with shared writeable memory如何使用共享可写内存并行化代码
【发布时间】:2017-06-06 11:23:21
【问题描述】:

我有一个目前非常慢且 CPU 很重的程序,我对如何并行化它感到困惑。

我的问题:

算法是这样工作的(使用随机数):

import numpy as np
import scipy.spatial as spatial
from collections import Counter

# Particle dictionary
num_parts = int(1e4)
particles = {'coords': np.random.rand(num_parts, 3),
             'flags': np.random.randint(0, 4, size=num_parts)}

# Object dictionary
num_objects = int(1e3)
objects_dict = {'x': [np.random.rand(1)
                      for i in range(0, num_objects)],
                'y': [np.random.rand(1)
                      for i in range(0, num_objects)],
                'z': [np.random.rand(1)
                      for i in range(0, num_objects)],
                'prop_to_calc': [np.array([])
                                 for i in range(0, num_objects)]}

# Build KDTree for particles
tree = spatial.cKDTree(particles['coords'])

# Loop through objects to calculate 'prop_to_calc' based
# on nearby particles flags within radius, r.
r = 0.1
for i in range(0, num_objects):

    # Find nearby indices of particles.
    indices = tree.query_ball_point([objects_dict['x'][i][0],
                                     objects_dict['y'][i][0],
                                     objects_dict['z'][i][0]],
                                    r)

    # Extract flags of nearby particles.
    flags_array = particles['flags'][indices]

    # Find most common flag and store that in property_to_calculate
    c = Counter(flags_array)
    objects_dict['prop_to_calc'] = np.append(objects_dict['prop_to_calc'],
                                             c.most_common(1)[0][0])

有两个数据集particlesobjects_dict。我想通过搜索半径内的附近粒子r 并找到它们最常见的标志来计算objects_dict['prop_to_dict']。这是通过cKDTreequery_ball_point 完成的。

对于这些数字,时间是: 10 loops, best of 3: 55.8 ms per loop 在 Ipython 4.2.0 中

但是,我想要num_parts=1e6num_objects=1e5,这会导致速度严重下降。

我的问题:

由于 CPU 很重,我想尝试并行化它以加快速度。

我查看了multiprocessingmulti threading。但是文档让我很困惑,我不确定如何将示例应用于此问题。

具体来说,我关心的是如何在进程之间共享这两个字典并在最后写入objects_dict

提前感谢您的帮助。

【问题讨论】:

  • 如果你想要速度,试试 python 编译器 Cython 或者只用 C(++) 编写你的程序。
  • 我考虑过 Cython。但是,我想先看看是否可以在我的示例中使用并行计算。

标签: python python-multithreading python-multiprocessing


【解决方案1】:

尝试拆分代码并制作本地列表和字典。联合最后的字典和列表。只需import threading 并创建一个带有子类threading.Threadclass CalculationThread 并在类的方法run 中编写并行代码,并通过thread = CalculationThread()thread.start() 启动线程

这个的最终代码是:

import threading

class CalculationThread(threading Tread):
    def __init__(self, range):
        self.range = range
        # use this copy instead for the thread
        self.objects_dict = dict(objects_dict)
    def run(self):
        for i in self.range:
            # your code

threads = []
threadCount = 8

lastThreadStart = num_objects
tLen = num_objects // threadCount
for i in reversed(range(threadCount)):
    threadStart = i*tLen
    t = CalculationThread(range(threadStart, lastThreadStart))
    lastThreadStart = threadStart

# wait for all to finish
for t in threads: t.join()

# code to union all objects_dicts comes here

【讨论】:

  • 这不会加快处理速度,实际上将执行拆分为线程实际上会减慢处理速度,因为现在不仅 Python 必须计算您的结果,但它还必须处理 GIL 和上下文切换。
  • 你没有任何意义。几乎每个独立的 Python 解释器(CPython、PyPy ...)都以不同的方式处理线程和进程——线程在同一个进程中执行,并且仅通过 GIL 相互隔离——这允许它们共享内存,但具有 GIL 和线程调度in place 实际上会减慢执行速度,因为解释器必须做额外的工作。当您的代码需要大量等待外部事件(例如 I/O)但对于计算,它会总是执行得比在单个线程中运行所有内容要慢。
  • 我认为,如果您将代码放入字符串中的函数中,执行它,字节码是每个线程单独生成的,所以它不应该是相同的字节码,然后它不应该受到影响吉尔。
  • 你想错了。除非您为每个代码 sn-p 启动子进程(在这种情况下为什么不首先使用 multiprocessing 模块),否则它将在同一个进程中运行。如果您启动不同的进程,您将无法共享内存(顺便说一句,这正是 multiprocessing 模块所做的)
  • @zwer 如果我理解正确,您可以使用docs 在进程池之间共享只读内存。你能分享每个进程也可以写的内存并在最后组合吗?
猜你喜欢
  • 2012-07-19
  • 1970-01-01
  • 2017-10-17
  • 2020-11-27
  • 2022-01-06
  • 1970-01-01
  • 1970-01-01
  • 2014-10-06
  • 2012-10-15
相关资源
最近更新 更多