【问题标题】:Python: Multiprocessing code is very slowPython:多处理代码非常慢
【发布时间】:2019-09-06 08:59:47
【问题描述】:

我使用pymongomongodb 一次性提取.8 百万条记录(这是一次性过程)并对其执行一些操作。

我的代码如下所示。

    proc = []
    for rec in cursor: # cursor has .8 million rows 
            print cnt
            cnt = cnt + 1
            url =  rec['urlk']
            mkptid = rec['mkptid']
            cii = rec['cii']

            #self.process_single_layer(url, mkptid, cii)


            proc = Process(target=self.process_single_layer, args=(url, mkptid, cii))
            procs.append(proc)
            proc.start()

             # complete the processes
    for proc in procs:
        proc.join()

process_single_layer 是一个基本上是从云端下载urls.并存储在本地的功能。

现在的问题是下载过程很慢,因为它必须点击一个 url。由于处理 1k 行的记录很大,因此需要 6 分钟。

为了减少我想实现Multiprocessing 的时间。但是很难看出上面的代码有什么不同。

请建议我如何在这种情况下提高性能。

【问题讨论】:

  • 我当然希望这不是你的程序的实际缩进,否则你在启动它后加入每个进程,有效地串联而不是并行运行它们。
  • 第二个不好的部分是您正在为每一行创建一个Process
  • 代码完全没有意义,也许你的识别是错误的?无论如何,我建议您在数据库中进行一些过滤,这样您就不会带来所有记录数据,而是您想要的。
  • 实际上我的多处理实现是不正确的。因为它不会增加任何价值,直到我们让它并行。但在这种情况下,我无法思考如何实现并行性
  • 首先,与其他建议一样,您需要在所有进程启动后加入,否则您的程序只是系列运行。其次,您是否要同时运行 .8m 进程。在具有修复连接条件的原始代码中,您的程序将生成 .8m 进程...

标签: python multiprocessing pymongo


【解决方案1】:

首先,您需要计算文件中的所有行数,然后生成固定数量的进程(理想情况下与处理器内核的数量相匹配),然后通过队列(每个进程一个)向这些进程提供一些行等于除法total_number_of_rows / number_of_cores。这种方法背后的想法是您在多个进程之间拆分这些行的处理,从而实现并行性。

一种动态找出核心数量的方法是:

import multiprocessing as mp
cores_count = mp.cpu_count()

通过创建队列列表循环添加行,然后在其上应用循环迭代器,可以通过避免初始行计数来实现轻微改进。

一个完整的例子:

import queue
import multiprocessing as mp
import itertools as itools

cores_count = mp.cpu_count()


def dosomething(q):

    while True:

        try:
            row = q.get(timeout=5)
        except queue.Empty:
            break

    # ..do some processing here with the row

    pass

if __name__ == '__main__':
    processes
    queues = []

    # spawn the processes
    for i in range(cores_count):
        q = mp.Queue()
        queues.append(q)
        proc = Process(target=dosomething, args=(q,))
        processes.append(proc)

    queues_cycle = itools.cycle(queues)
    for row in cursor:
        q = next(queues_cycle)
        q.put(row)

    # do the join after spawning all the processes
    for p in processes:
        p.join()

【讨论】:

  • 如果处理受 CPU 限制,这是正确的,但似乎 OP 的问题是他正在从 Internet 获取 URL - 这将受网络限制。他可以承受比 CPU 多 许多 个线程,因为大多数线程将被阻塞等待网络响应。
  • OP 最好使用某种基于消息的系统,这样不会阻塞线程。在这一点上,他最好还是保持单线程(因为多线程很难理解。
  • @MartinBonner "...因为多线程很难理解..." 并不是一个真正的论点。使用线程而不是多处理对于网络操作来说并不是一个好主意,因为为什么线程在单个进程中运行并且并行化是虚拟的(它们实际上是轮流从每个线程运行每个操作),这意味着不是纯粹的并行性。我已经使用网络操作完成了多处理,老实说,与线程或简单线程相比,我得到了相当不错的结果。
【解决方案2】:

在这种情况下使用池更容易。

队列不是必需的,因为您不需要在生成的进程之间进行通信。我们可以使用Pool.map 来分配工作负载。

Pool.imapPool.imap_unordered 可能会在更大的块大小下更快。 (参考:https://docs.python.org/3/library/multiprocessing.html#multiprocessing.pool.Pool.imap)如果需要,您可以使用 Pool.starmap 并摆脱元组拆包。

from multiprocessing import Pool

def process_single_layer(data):
    # unpack the tuple and do the processing
    url, mkptid, cii = data
    return "downloaded" + url

def get_urls():
    # replace this code: iterate over cursor and yield necessary data as a tuple
    for rec in range(8): 
            url =  "url:" + str(rec)
            mkptid = "mkptid:" + str(rec)
            cii = "cii:" + str(rec)
            yield (url, mkptid, cii)

#  you can come up with suitable process count based on the number of CPUs.
with Pool(processes=4) as pool:
    print(pool.map(process_single_layer, get_urls()))

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-11-11
    • 1970-01-01
    • 2020-03-02
    • 1970-01-01
    • 2022-06-13
    • 2019-10-15
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多