【发布时间】:2015-09-30 10:33:05
【问题描述】:
我最近发现了dask 模块,旨在成为一个易于使用的python 并行处理模块。对我来说最大的卖点是它适用于 pandas。
在阅读了它的手册页之后,我找不到一种方法来完成这个琐碎的并行化任务:
ts.apply(func) # for pandas series
df.apply(func, axis = 1) # for pandas DF row apply
目前,为了实现这一目标,AFAIK,
ddf.assign(A=lambda df: df.apply(func, axis=1)).compute() # dask DataFrame
这是一种丑陋的语法,实际上比完全慢
df.apply(func, axis = 1) # for pandas DF row apply
有什么建议吗?
编辑:感谢@MRocklin 提供地图功能。它似乎比普通熊猫应用要慢。这与熊猫 GIL 发布问题有关还是我做错了?
import dask.dataframe as dd
s = pd.Series([10000]*120)
ds = dd.from_pandas(s, npartitions = 3)
def slow_func(k):
A = np.random.normal(size = k) # k = 10000
s = 0
for a in A:
if a > 0:
s += 1
else:
s -= 1
return s
s.apply(slow_func) # 0.43 sec
ds.map(slow_func).compute() # 2.04 sec
【问题讨论】:
-
我对@987654329@模块不熟悉。对于多重处理,当我必须逐行处理大数据帧时,python 模块
multiprocessing非常适合我。思路也很简单:使用np.array_split将大数据框拆分为8个,使用multiprocessing同时处理;完成后,使用pd.concat将它们连接回原始长度。有关完整代码示例的相关帖子,请参阅stackoverflow.com/questions/30904354/… -
谢谢,非常好。多处理模块的问题是您需要有一个命名函数(不是 lambda)并将其放在 name=="main" 块之外。这使得研究代码的结构很糟糕。
-
如果您只想使用更好的多处理,您可以查看@mike-mckerns 的multiprocess。您也可以尝试使用 dask core 而不是 dask.dataframe 并构建字典或使用类似 github.com/ContinuumIO/dask/pull/408
标签: python pandas parallel-processing dask