【问题标题】:Parallelising a plotting loop in Jupyter notebook在 Jupyter 笔记本中并行化绘图循环
【发布时间】:2022-08-07 12:24:36
【问题描述】:

我正在使用 Python 3.5.1 版。我想并行化一个循环,该循环用于使用 imshow 绘制一组数组。没有任何并行化的最小代码如下

import matplotlib.pyplot as plt
import numpy as np

# Generate data

arrays   = [np.random.rand(3,2) for x in range(10)]
arrays_2 = [np.random.rand(3,2) for x in range(10)]

# Loop and plot sequentially

for i in range(len(arrays)):

    # Plot side by side

    figure = plt.figure(figsize = (20, 12))
    ax_1 = figure.add_subplot(1, 2, 1)
    ax_2 = figure.add_subplot(1, 2, 2)

    ax_1.imshow(arrays[i], interpolation=\'gaussian\', cmap=\'RdBu\', vmin=0.5*np.min(arrays[i]), vmax=0.5*np.max(arrays[i]))
    ax_2.imshow(arrays_2[i], interpolation=\'gaussian\', cmap=\'YlGn\', vmin=0.5*np.min(arrays_2[i]), vmax=0.5*np.max(arrays_2[i]))

    plt.savefig(\'./Figure_{}\'.format(i), bbox_inches=\'tight\')
    plt.close()

此代码目前是在 Jupyter 笔记本中编写的,我只想通过 Jupyter 笔记本进行所有处理。虽然这很好用,但实际上我有 2500 多个阵列,并且每秒大约 1 个绘图,这需要很长时间才能完成。我想做的是在 N 个处理器之间拆分计算,以便每个处理器为 len(arrays)/N 个数组绘制图。由于这些图是各个阵列本身的图,因此在任何计算过程中,内核无需相互通信(无共享)。

我已经看到multiprocessing package 对类似问题有好处。但是,它不适用于我的问题,因为您不能将二维数组传递给函数。如果我这样修改上面的代码

# Generate data

arrays   = [np.random.rand(3,2) for x in range(10)]
arrays_2 = [np.random.rand(3,2) for x in range(10)]

x = list(zip(arrays, arrays_2))

def plot_file(information):

    arrays, arrays_2 = list(information[0]), list(information[1])
    print(np.shape(arrays[0][0]), np.shape(arrays_2[0][0]))
    
    # Loop and plot sequentially

    for i in range(len(arrays)):        

        # Plot side by side

        figure = plt.figure(figsize = (20, 12))
        ax_1 = figure.add_subplot(1, 2, 1)
        ax_2 = figure.add_subplot(1, 2, 2)

        ax_1.imshow(arrays[i], interpolation=\'gaussian\', cmap=\'RdBu\', vmin=0.5*np.min(arrays[i]), vmax=0.5*np.max(arrays[i]))
        ax_2.imshow(arrays_2[i], interpolation=\'gaussian\', cmap=\'YlGn\', vmin=0.5*np.min(arrays_2[i]), vmax=0.5*np.max(arrays_2[i]))

        plt.savefig(\'./Figure_{}\'.format(i), bbox_inches=\'tight\')
        plt.close()
    
from multiprocessing import Pool
pool = Pool(4)
pool.map(plot_file, x)

然后我得到错误 \'TypeError: Invalid dimensions for image data\' 并且数组维度的打印输出现在只是 (2, ) 而不是 (3, 2)。显然,这是因为多处理不能/不能将 2D 数组作为输入处理。

所以我想知道,如何在 Jupyter 笔记本中并行化?有人可以告诉我怎么做吗?

  • 这回答了你的问题了吗? How do I parallelize a simple Python loop? 使用multiprocessing.Pool 记下答案。
  • 一个问题 - 为什么不在每个函数内生成/准备数组,而不是提前?
  • @MichaelDelgado当我在函数内生成数据时,上面的多处理代码有效。但是,如果我使用 Pool(4) 运行代码,那么我很确定每个处理器只是在整个数组集上进行计算,并且数据不会在四个处理器之间均匀分布,因为代码占用的数量完全相同在没有多处理的情况下计算时间。我想要的是将 N 个处理器中的数据均匀地拆分为 N 个子集,并让单个处理器仅在数组的单个子集上进行计算。
  • 对...所以不要让每个处理器都处理全套作业。或者您可以设置更多的工作模型,并让它们都使用队列中的任务。
  • 是的,不,您需要明确说明任务的分配方式。您可以使用 multiprocessing.map,类似于我在回答中调用 dask 的方式。您是否有不想使用 dask 的原因?这是一个很棒的包:)

标签: python parallel-processing


【解决方案1】:

一种简单的方法是使用多处理引擎使用dask.distributed。我只建议一个外部模块,因为 dask 会为您处理对象的序列化,这使得这成为一个非常简单的操作:

import matplotlib
# include this line to allow your processes to plot without a screen
matplotlib.use('Agg')

import matplotlib.pyplot as plt
import dask.distributed
import numpy as np

def plot_file(i, array_1, array_2):
    matplotlib.use('Agg')

    # will be called once for each array "job"
    figure = plt.figure(figsize = (20, 12))
    ax_1 = figure.add_subplot(1, 2, 1)
    ax_2 = figure.add_subplot(1, 2, 2)

    for ax, arr, cmap in [(ax_1, array_1, 'RdBu'), (ax_2, array_2, 'YlGn')]:
        ax.imshow(
            arr,
            interpolation='gaussian',
            cmap='RdBu',
            vmin=0.5*np.min(arr),
            vmax=0.5*np.max(arr),
        )

    figure.savefig('./Figure_{}'.format(i), bbox_inches='tight')
    plt.close(figure)

arrays   = [np.random.rand(3,2) for x in range(10)]
arrays_2 = [np.random.rand(3,2) for x in range(10)]

client = dask.distributed.Client() # uses multiprocessing by default
futures = client.map(plot_file, range(len(arrays)), arrays, arrays_2)
dask.distributed.progress(futures)

但是,如果可能,更有效的是在映射任务中生成或准备数组。这也将允许您并行执行数组操作、I/O 等:

def prep_arrays_and_plot(i):
    array_1 = np.random.rand(3,2)
    array_2 = np.random.rand(3,2)
    plot_file(i, array_1, array_2)

futures = client.map(prep_arrays_and_plot, range(10))
dask.distributed.progress(futures)

在这一点上,您不需要腌制任何东西,因此使用多处理编写并不是什么大不了的事。以下脚本运行良好:

import matplotlib
matplotlib.use("Agg")

import matplotlib.pyplot as plt
import numpy as np
import multiprocessing

def plot_file(i, array_1, array_2):
    matplotlib.use('Agg')

    # will be called once for each array "job"
    figure = plt.figure(figsize = (20, 12))
    ax_1 = figure.add_subplot(1, 2, 1)
    ax_2 = figure.add_subplot(1, 2, 2)

    for ax, arr, cmap in [(ax_1, array_1, 'RdBu'), (ax_2, array_2, 'YlGn')]:
        ax.imshow(
            arr,
            interpolation='gaussian',
            cmap='RdBu',
            vmin=0.5*np.min(arr),
            vmax=0.5*np.max(arr),
        )

    figure.savefig('./Figure_{}'.format(i), bbox_inches='tight')
    plt.close(figure)

def prep_arrays_and_plot(i):
    array_1 = np.random.rand(3,2)
    array_2 = np.random.rand(3,2)
    plot_file(i, array_1, array_2)

def main():
    pool = multiprocessing.Pool(4)
    pool.map(prep_arrays_and_plot, range(10))

if __name__ == "__main__":
    main()

请注意,如果您从 jupyter 笔记本运行此程序,则不能简单地在单元格中定义函数并将它们传递给 multiprocessing.Pool。相反,您必须在不同的文件中定义它们并导入它们。这不适用于 dask(实际上,如果您在 notebook 中使用 dask 定义函数会更容易)。

【讨论】:

    猜你喜欢
    • 2022-07-21
    • 2015-10-15
    • 2019-03-17
    • 2021-07-27
    • 2021-03-22
    • 2016-07-04
    • 2020-07-22
    • 2016-04-11
    • 1970-01-01
    相关资源
    最近更新 更多