【问题标题】:Giving a large dictionary to a async function is making the code very slow为异步函数提供大字典会使代码非常慢
【发布时间】:2013-03-11 21:37:01
【问题描述】:

我在我的 python 代码中使用多处理来异步运行一个函数:

import multiprocessing

po = multiprocessing.Pool()
for elements in a_list:
    results.append(po.apply_async(my_module.my_function, (some_arguments, elements, a_big_argument)))               
po.close()
po.join()
for r in results:
    a_new_list.add(r.get())

a_big_argument 是一个字典。我把它作为一个论据。从某种意义上说,它在 10 到 100 个月之间很大。它似乎对我的代码的性能有很大的影响。

我可能在这里做了一些愚蠢且效率不高的事情,因为我的代码的性能确实因这个新参数而下降。

处理大字典的最佳方法是什么?我不想每次都在我的函数中加载它。是否可以改为创建数据库并连接到它?

这是您可以运行的代码:

'''
Created on Mar 11, 2013

@author: Antonin
'''

import multiprocessing
import random

# generate an artificially big dictionary
def generateBigDict():
    myBigDict = {}
    for key in range (0,1000000):
        myBigDict[key] = 1
    return myBigDict

def myMainFunction():
    # load the dictionary
    myBigDict = generateBigDict()
    # create a list on which we will asynchronously run the subfunction
    myList = []
    for list_element in range(0,20):
        myList.append(random.randrange(0,1000000))
    # an empty set to receive results
    set_of_results = set()
    # there is a for loop here on one of the arguments
    for loop_element in range(0,150):
        results = []
        # asynchronoulsy run the subfunction
        po = multiprocessing.Pool()
        for list_element in myList:
            results.append(po.apply_async(mySubFunction, (loop_element, list_element, myBigDict)))               
        po.close()
        po.join()
        for r in results:
            set_of_results.add(r.get())
    for element in set_of_results:
        print element

def mySubFunction(loop_element, list_element, myBigDict):
    import math
    intermediaryResult = myBigDict[list_element]
    finalResult = intermediaryResult + loop_element
    return math.log(finalResult)

if __name__ == '__main__':
    myMainFunction()

【问题讨论】:

  • 你能提供一个小的真实程序来演示这个问题吗?见SSCCE.org
  • 你在windows上吗?在 Windows 上 multiprocessing 必须 腌制参数然后发送到子进程,而在 unix 上 可能 能够fork(这应该更有效)。
  • 我在我的 linux 机器上进行了测试,通过包含大约 200k 个对象的 dict 大约需要 0.05 秒。
  • 我刚刚添加了一个(简单?小?现实?)真实程序。我在 UNIX 上工作。需要探索fork,不知道怎么用。

标签: python asynchronous multiprocessing


【解决方案1】:

我使用multiprocessing.Manager 来做到这一点。

import multiprocessing

manager = multiprocessing.Manager()
a_shared_big_dictionary = manager.dict(a_big_dictionary)

po = multiprocessing.Pool()
for elements in a_list:
    results.append(po.apply_async(my_module.my_function, (some_arguments, elements, a_shared_big_dictionary)))               
po.close()
po.join()
for r in results:
    a_new_list.add(r.get())

现在,它更快了。

【讨论】:

    【解决方案2】:

    查看Shared-memory objects in python multiprocessing问题的答案。

    它建议使用multiprocessing.Array 将数组传递给子进程或使用fork()。

    【讨论】:

    • 感谢您的回答。阅读您的链接非常有帮助。知道我想在我的进程之间共享的对象是一个字典(而不是一个数组),使用Manager 不是更好吗?
    • 您正在寻找速度提升。管理器在将 Python 对象传递给子进程时涉及序列化(酸洗)和反序列化(如果我理解正确的话)。这应该比 Array 等共享内存解决方案慢。不过,请自行检查您获得的性能。
    • 查看您的解决方案(使用 Manager),我发现您获得了所需的性能提升。好:)
    【解决方案3】:

    您传递给Pool 方法之一的任何参数(例如apply_async)都需要被腌制,通过管道发送到工作进程,并在工作进程中取消腌制。这个 pickle/pass/unpickle 过程在时间和内存上可能会很昂贵,特别是如果您有一个大型对象图,因为每个工作进程都必须创建一个单独的副本。

    根据问题的具体形式,有许多不同的方法可以避免这些泡菜。由于您的工作人员只是阅读您的字典而不是写入它,因此您可以安全地直接从您的函数中引用它(即不将其传递给apply_async)并依靠fork() 来避免创建一个在工作进程中复制。

    更好的是,您可以更改mySubFunction(),使其接受intermediaryResult 作为参数,而不是使用list_elementmyBigDict 查找它。 (您可以通过闭包来做到这一点,但我不能 100% 确定 pickle 也不会尝试复制封闭的 myBigDict 对象。)

    或者,您可以将myBigDict 放在所有进程都可以安全共享它的地方,例如one of the simple persistance methods, such as dbm or sqlite,并让工作人员从那里访问它。

    不幸的是,所有这些解决方案都要求您更改任务函数的形状。避免这种“变形”是人们喜欢“真正的”cpu 线程的原因之一。

    【讨论】:

    • 感谢您的回答。您提到“所有进程都可以安全共享它的某个地方”。你觉得Server process with multiprocessing.Manager会对应这个吗?
    • 经理可能也可以工作,但请记住,它只提供一致性,而不是原子性、完整性或持久性——即如果你共享一个数据结构来写,你仍然需要自己管理并发原语(例如使用共享锁——管理器也可以为你持有)。
    • 在您的情况下,无论如何都不需要管理器和代理对象。正如我所说,您可以将大字典设置为工作方法并直接访问它,并依靠fork() 写入时复制行为来确保它不会被复制。
    • 我发现要找到有关fork() 的信息要困难得多。你建议使用它,但我不知道该怎么做。经理非常易于使用,一行可能需要 3 个单词,我的功能没有变化。它在 2 分钟内完成,我确切地知道它的作用。如果您对fork() 有任何参考,我会很感兴趣。
    • 您不必做任何事情(除了在创建工人之前创建数据结构——您已经这样做了)。 multiprocessing 已经使用 fork() 来创建工人。你没有明确使用 fork(),你只是依赖它是copy-on-write进程产生行为....
    猜你喜欢
    • 2019-03-09
    • 1970-01-01
    • 2022-01-17
    • 2023-03-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-03-19
    • 1970-01-01
    相关资源
    最近更新 更多