【问题标题】:Python parallelised correlation slower than single process correlationPython并行相关比单进程相关慢
【发布时间】:2017-02-16 09:30:27
【问题描述】:

我想在 Python 中使用 multiprocessing 模块并行化 df.corr()。我正在取一列并计算相关值与一个进程中的所有列和第二列与另一个进程中的其他列。我将继续以这种方式通过堆叠来自所有进程的结果行来填充相关矩阵的上层。

我取了形状(678461, 210)的样本数据,并尝试了我的并行化方法和df.corr(),分别得到了214.40s42.64s的运行时间。所以,我的并行化方法需要更多时间。

有什么办法可以改善吗?

import multiprocessing as mp
import pandas as pd
import numpy as np
from time import *

def _correlation(args):

    i, mat, mask = args
    ac = mat[i]

    arr = []

    for j in range(len(mat)):  
        if i > j:
            continue

        bc = mat[j]
        valid = mask[i] & mask[j]
        if valid.sum() < 1:
            c = NA    
        elif i == j:
            c = 1.
        elif not valid.all():
            c = np.corrcoef(ac[valid], bc[valid])[0, 1]
        else:
            c = np.corrcoef(ac, bc)[0, 1]

        arr.append((j, c))

    return arr

def correlation_multi(df):

    numeric_df = df._get_numeric_data()
    cols = numeric_df.columns
    mat = numeric_df.values

    mat = pd.core.common._ensure_float64(mat).T
    K = len(cols)
    correl = np.empty((K, K), dtype=float)
    mask = np.isfinite(mat)

    pool = mp.Pool(processes=4)

    ret_list = pool.map(_correlation, [(i, mat, mask) for i in range(len(mat))])

    for i, arr in enumerate(ret_list):
        for l in arr:
            j = l[0]
            c = l[1]

            correl[i, j] = c
            correl[j, i] = c

    return pd.DataFrame(correl, index = cols, columns = cols)

if __name__ == '__main__':
    noise  = pd.DataFrame(np.random.randint(0,100,size=(100000, 50)))
    noise2  = pd.DataFrame(np.random.randint(100,200,size=(100000, 50)))
    df = pd.concat([noise, noise2], axis=1)

    #Single process correlation    
    start = time()
    s = df.corr()
    print('Time taken: ',time()-start)

    #Multi process correlation
    start = time()
    s1 = correlation_multi(df)
    print('Time taken: ',time()-start)

【问题讨论】:

    标签: python pandas parallel-processing multiprocessing correlation


    【解决方案1】:

    _correlation 的结果必须通过进程间通信从工作进程转移到运行Pool 的进程。

    这意味着返回的数据被腌制,发送到其他进程,取消腌制并添加到结果列表中。 这需要时间,而且本质上是一个连续的过程。

    map 按照发送顺序处理退货,IIRC。因此,如果一次迭代花费的时间相对较长,则其他结果可能会一直等待。您可以尝试使用imap_unordered,它会在它们到达后立即产生结果。

    【讨论】:

      猜你喜欢
      • 2013-08-04
      • 2019-05-18
      • 1970-01-01
      • 1970-01-01
      • 2019-01-26
      • 2013-11-17
      • 2013-06-08
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多