【问题标题】:how to output results of python parallel computing (ipython-parallel or multiprocessing) to a pandas dataframe?如何将 python 并行计算(ipython-parallel 或 multiprocessing)的结果输出到 pandas 数据帧?
【发布时间】:2015-08-03 13:45:27
【问题描述】:

简单的问题:我读过的所有教程都向您展示了如何使用 ipython.parallel 或多处理将并行计算的结果输出到列表(或最好是字典)。

您能否指出一个使用任一库将计算结果输出到共享 pandas 数据框的简单示例?

http://gouthamanbalaraman.com/blog/distributed-processing-pandas.html - 本教程向您展示如何读取输入数据帧(下面的代码),但是我将如何将 4 个并行计算的结果输出到一个数据帧?

import pandas as pd
import multiprocessing as mp

LARGE_FILE = "D:\\my_large_file.txt"
CHUNKSIZE = 100000 # processing 100,000 rows at a time

def process_frame(df):
        # process data frame
        return len(df)

if __name__ == '__main__':
        reader = pd.read_table(LARGE_FILE, chunksize=CHUNKSIZE)
        pool = mp.Pool(4) # use 4 processes

        funclist = []
        for df in reader:
                # process each data frame
                f = pool.apply_async(process_frame,[df])
                funclist.append(f)

        result = 0
        for f in funclist:
                result += f.get(timeout=10) # timeout in 10 seconds

        print "There are %d rows of data"%(result)

【问题讨论】:

  • 你为什么不把你的输出放到一个列表中,然后把它减少到一个数据框,比如reduce(lambda x,y: x.append(y), your_list)
  • 您需要向我们展示您正在尝试做什么,向我们展示您的单线程解决方案以及您打算如何对其进行多处理。

标签: python pandas parallel-processing multiprocessing ipython-parallel


【解决方案1】:

您要求multiprocessing(或其他python 并行模块)输出到它们不直接输出到的数据结构。如果你使用来自任何并行包的Pool,你最好得到一个列表(使用map)或一个迭代器(使用imap)。如果您使用来自multiprocessing 的共享内存,您可能能够将结果放入一个内存块中,该内存块可以通过ctypes 的指针访问。

那么问题是,您能否将结果从迭代器或共享内存块提取到pandas.DataFrame 中?我认为答案是肯定的。是的你可以。但是,我认为我没有在教程中看到过这样做的简单示例……因为它并不那么简单。

迭代器路由似乎不太可能,因为您需要获取numpy 来消化迭代器,而不会将结果作为列表首先拉回python。我会选择共享内存路线。我认为这应该给你一个DataFrame 的输出,然后你可以在multiprocessing 中使用它:

from multiprocessing import sharedctypes as sh
from numpy import ctypeslib as ct        
import pandas as pd

ra = sh.RawArray('i', 4)
arr = ct.as_array(ra)
arr.shape = (2,2)
x = pd.DataFrame(arr)

那么您所要做的就是将数组的句柄传递给multiprocessing.Process

import multiprocessing as mp
p1 = mp.Process(target=doit, args=(arr[:1, :], 1))
p2 = mp.Process(target=doit, args=(arr[1:, :], 2))
p1.start()
p2.start()
p1.join()
p2.join()

然后,通过一些指针魔术,结果应该填入你的DataFrame .

我会让你编写 doit 函数来随心所欲地操作数组。

编辑:这看起来像是使用类似方法的好答案...https://stackoverflow.com/a/22487898/2379433。这似乎也有效:https://stackoverflow.com/a/27027632/2379433

【讨论】:

  • 非常感谢您的回答!
猜你喜欢
  • 2017-12-10
  • 2021-04-15
  • 2015-10-11
  • 2015-06-27
  • 1970-01-01
  • 2013-06-10
  • 1970-01-01
  • 2020-12-07
  • 2014-03-05
相关资源
最近更新 更多