【问题标题】:How to use Multiprocessing Pool on a dictionary input in Python?如何在 Python 中的字典输入上使用多处理池?
【发布时间】:2023-02-23 07:12:55
【问题描述】:

我的目标是在 python 字典上使用 map reduce aka Multiprocessing pool。我希望它将键值对映射到不同的核心,然后将结果聚合为字典。

from multiprocessing.pool import Pool
elems = {i:i for i in range(1_000_000)}

def func(x):
    return (x, elems[x]**2)

with Pool() as pool:
    results = pool.map(func, elems.keys())
    results = {a:b for a,b in results}

这是一个有点棘手的解决方案,但是否有更 Pythonic 的方式来接收字典输入并使用 Python 中的多处理池生成字典输出?

【问题讨论】:

  • 不清楚你的意思。什么输入? pool.map 的输入可能只是elems(实际上等同于elems.keys())...所以从这个意义上说,输入是dict。那么你到底想要什么?我不清楚。如果要映射键值对,那么使用elems.items(),那么x就是一个键值对。
  • 我假设pool.map按顺序返回结果,如果是这样,为什么不只做results = dict(zip(elems.keys(), results)),让results只返回elems[x]**2
  • 顺便说一句,results = {a:b for a,b in results}只能是results = dict(results),一般来说,{k:v for k,v in whatever}只能是dict(whatever)

标签: python dictionary multiprocessing


【解决方案1】:

您可以使用 ProcessPoolExecutor 轻松进行 map-reduce:

from concurrent.futures import ProcessPoolExecutor


def process(item):
    return (item[0], item[1] ** 2)


def main():
    elems = {i: i for i in range(1_000_000)}

    output = {}

    with ProcessPoolExecutor() as pool:
        results = pool.map(process, elems.items(), chunksize=1_000)

        for result in results:
            output[result[0]] = result[1]

    print(output)


if __name__ == "__main__":
    main()

在这里,results 是一个(异步)迭代器,每次并行处理的结果可用时它都会产生一个值,您可以对其进行迭代以减少部分。由于数据的进程间通信,多处理的成本可能很高,因此,您应该调整 chunksize 参数以适合您的用例。

多处理的一些建议:

  • 永远不要改变将要并发执行的函数的共享状态,这会导致数据竞争;
  • if __name == "__main__" 中保护您的主要功能,否则您会遇到问题;
  • 不要将大数据声明为全局状态(例如你的elems字典),否则它将在每个子进程中被复制,声明一个主函数;
  • 在并发执行的函数中尽可能避免访问共享状态,即使用def func(key, value)而不是elems[x]

【讨论】:

  • 选择块大小的好方法是什么?假设我有 3000 个元素和 30 个核心,100 个是否是一个不错的选择,因为这会在进程之间平均分配任务?或者也许我想得太多了。
  • 这取决于内存中元素的大小。一般来说,只需尝试不同的值并检查程序的速度和内存消耗即可估计好的价值。无论如何,你的建议似乎是最佳的。
  • 我在回答中添加了一些细节。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-04-26
  • 2017-10-12
  • 2013-01-26
  • 2023-02-22
  • 1970-01-01
  • 2020-03-23
  • 2020-08-15
相关资源
最近更新 更多