【发布时间】:2015-01-03 05:40:44
【问题描述】:
我正在尝试对 pandas 数据帧使用多处理,即将数据帧拆分为 8 个部分。使用 apply 对每个部分应用一些功能(每个部分在不同的过程中处理)。
编辑: 这是我终于找到的解决方案:
import multiprocessing as mp
import pandas.util.testing as pdt
def process_apply(x):
# do some stuff to data here
def process(df):
res = df.apply(process_apply, axis=1)
return res
if __name__ == '__main__':
p = mp.Pool(processes=8)
split_dfs = np.array_split(big_df,8)
pool_results = p.map(aoi_proc, split_dfs)
p.close()
p.join()
# merging parts processed by different processes
parts = pd.concat(pool_results, axis=0)
# merging newly calculated parts to big_df
big_df = pd.concat([big_df, parts], axis=1)
# checking if the dfs were merged correctly
pdt.assert_series_equal(parts['id'], big_df['id'])
【问题讨论】:
-
@yemu 你到底想通过这段代码实现什么?
-
目前只应用饱和 CPU 的一个核心。我想使用多进程并使用所有内核来减少处理时间
-
如果你把问题放在一边,然后把答案放在答案中会更好。这样我们就可以在不查看变更日志的情况下看到更多的过程。
-
“aoi_proc”应该是“进程”吗?也许将您的“进程”函数重命名为简单的“f”在多处理上下文中会更具可读性
-
我对 process_apply 应该是什么样子感到困惑。我的函数是行的函数。比如:
def process_apply(rw): return(rw['A']*rw['B'])。对吗?
标签: python pandas multiprocessing