【问题标题】:How to parallelize iteration over a range, using StdLib and Python 3?如何使用 StdLib 和 Python 3 在一个范围内并行化迭代?
【发布时间】:2018-10-03 21:34:56
【问题描述】:

我这几天一直在寻找答案,但无济于事。我可能只是不理解那里漂浮的部分,multiprocessing 模块上的 Python 文档相当大,我不清楚。

假设你有以下 for 循环:

import timeit


numbers = []

start = timeit.default_timer()

for num in range(100000000):
    numbers.append(num)

end = timeit.default_timer()

print('TIME: {} seconds'.format(end - start))
print('SUM:', sum(numbers))

输出:

TIME: 23.965870224497916 seconds
SUM: 4999999950000000

对于这个例子,假设你有一个 4 核处理器。有没有办法一共创建 4 个进程,每个进程都在一个单独的 CPU 内核上运行并且完成速度大约快 4 倍,所以 24s/4 个进程 = ~6 秒?

以某种方式将 for 循环分成 4 个相等的块,然后将这 4 个块添加到 numbers 列表中以等于相同的总和?有这个 stackoverflow 线程:Parallel Simple For Loop 但我不明白。谢谢大家。

【问题讨论】:

    标签: python python-3.x parallel-processing multiprocessing range


    【解决方案1】:

    是的,这是可行的。您的计算不依赖于中间结果,因此您可以轻松地将任务分成块并将其分配到多个进程中。这就是所谓的

    尴尬的并行问题

    这里唯一棘手的部分可能是,首先将范围分成相当相等的部分。直出我个人的 lib 两个函数来处理这个问题:

    # mp_utils.py
    
    from itertools import accumulate
    
    def calc_batch_sizes(n_tasks: int, n_workers: int) -> list:
        """Divide `n_tasks` optimally between n_workers to get batch_sizes.
    
        Guarantees batch sizes won't differ for more than 1.
    
        Example:
        # >>>calc_batch_sizes(23, 4)
        # Out: [6, 6, 6, 5]
    
        In case you're going to use numpy anyway, use np.array_split:
        [len(a) for a in np.array_split(np.arange(23), 4)]
        # Out: [6, 6, 6, 5]
        """
        x = int(n_tasks / n_workers)
        y = n_tasks % n_workers
        batch_sizes = [x + (y > 0)] * y + [x] * (n_workers - y)
    
        return batch_sizes
    
    
    def build_batch_ranges(batch_sizes: list) -> list:
        """Build batch_ranges from list of batch_sizes.
    
        Example:
        # batch_sizes [6, 6, 6, 5]
        # >>>build_batch_ranges(batch_sizes)
        # Out: [range(0, 6), range(6, 12), range(12, 18), range(18, 23)]
        """
        upper_bounds = [*accumulate(batch_sizes)]
        lower_bounds = [0] + upper_bounds[:-1]
        batch_ranges = [range(l, u) for l, u in zip(lower_bounds, upper_bounds)]
    
        return batch_ranges
    

    那么您的主脚本将如下所示:

    import time
    from multiprocessing import Pool
    from mp_utils import calc_batch_sizes, build_batch_ranges
    
    
    def target_foo(batch_range):
        return sum(batch_range)  # ~ 6x faster than target_foo1
    
    
    def target_foo1(batch_range):
        numbers = []
        for num in batch_range:
            numbers.append(num)
        return sum(numbers)
    
    
    if __name__ == '__main__':
    
        N = 100000000
        N_CORES = 4
    
        batch_sizes = calc_batch_sizes(N, n_workers=N_CORES)
        batch_ranges = build_batch_ranges(batch_sizes)
    
        start = time.perf_counter()
        with Pool(N_CORES) as pool:
            result = pool.map(target_foo, batch_ranges)
            r_sum = sum(result)
        print(r_sum)
        print(f'elapsed: {time.perf_counter() - start:.2f} s')
    

    请注意,我还将您的 for 循环切换为对 range 对象进行简单求和,因为它提供了更好的性能。如果您无法在实际应用中执行此操作,列表理解仍然比示例中手动填充列表快约 60%。

    示例输出:

    4999999950000000
    elapsed: 0.51 s
    
    Process finished with exit code 0
    

    【讨论】:

    • @probat 是的,抱歉忘记了
    • 你必须有一台超级计算机获得 0.51 秒 :)
    • @probat 6 岁机器 ;) 但我在 Linux 上。当你在 Windows 上时,启动进程需要更长的时间,因为 Windows 没有分叉。
    【解决方案2】:
    import timeit
    
    from multiprocessing import Pool
    
    def appendNumber(x):
        return x
    
    start = timeit.default_timer()
    
    with Pool(4) as p:
        numbers = p.map(appendNumber, range(100000000))
    
    end = timeit.default_timer()
    
    print('TIME: {} seconds'.format(end - start))
    print('SUM:', sum(numbers))
    

    所以Pool.map 就像内置的map 函数。它接受一个函数和一个可迭代对象,并生成一个在可迭代对象的每个元素上调用该函数的结果列表。在这里,由于我们实际上并不想更改可迭代范围内的元素,因此我们只返回参数。

    关键是Pool.map 将提供的可迭代对象(此处为range(1000000000))分成块并将它们发送到它拥有的进程数(此处定义为Pool(4) 中的4),然后将结果重新加入一份清单。

    运行时我得到的输出是

    TIME: 8.748245699999984 seconds
    SUM: 4999999950000000
    

    【讨论】:

      【解决方案3】:

      我做了一个对比,有时候任务拆分的时间可能会更长:

      文件multiprocessing_summation.py

      def summation(lst):
        sum = 0
        for x in range(lst[0], lst[1]):
          sum += x
        return sum
      

      文件multiprocessing_summation_master.py

      %%file ./examples/multiprocessing_summation_master.py
      import multiprocessing as mp
      import timeit
      import os
      import sys
      import multiprocessing_summation as mps
      
      if __name__ == "__main__":
      
        if len(sys.argv) == 1:
          print(f'{sys.argv[0]} <number1 ...>')
          sys.exit(1)
        else:
          args = [int(x) for x in sys.argv[1:]]
      
        nBegin = 1
        nCore = os.cpu_count()
      
        for nEnd in args:
      
          ### Approach 1  ####
          ####################
          start = timeit.default_timer()
          answer1 = mps.summation((nBegin, nEnd+1))
          end = timeit.default_timer()
          print(f'Answer1 = {answer1}')
          print(f'Time taken = {end - start}')
      
          ### Approach 2 ####
          ####################
          start = timeit.default_timer()
          lst = []
          for x in range(nBegin, nEnd, int((nEnd-nBegin+1)/nCore)):
            lst.append(x)
          lst.append(nEnd+1)
      
          lst2 = []
          for x in range(1, len(lst)):
            lst2.append((lst[x-1], lst[x]))
      
          with mp.Pool(processes=nCore) as pool:
            answer2 = pool.map(mps.summation, lst2)
          end = timeit.default_timer()
          print(f'Answer2 = {sum(answer2)}')
          print(f'Time taken = {end - start}')
      

      运行第二个脚本:

      python multiprocessing_summation_master.py 1000 100000 10000000 1000000000

      输出是:

      Answer1 = 500500
      Time taken = 4.558405389566795e-05
      Answer2 = 500500
      Time taken = 0.15728066685459452
      Answer1 = 5000050000
      Time taken = 0.005781152051264199
      Answer2 = 5000050000
      Time taken = 0.14532123447452705
      Answer1 = 50000005000000
      Time taken = 0.4903863230334036
      Answer2 = 50000005000000
      Time taken = 0.49744346392131533
      Answer1 = 500000000500000000
      Time taken = 50.825169837068
      Answer2 = 500000000500000000
      Time taken = 26.603663061636567
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2012-11-05
        • 2011-10-28
        • 2012-11-09
        • 2011-01-28
        • 1970-01-01
        • 1970-01-01
        • 2015-10-27
        • 1970-01-01
        相关资源
        最近更新 更多