【发布时间】:2018-01-14 15:41:09
【问题描述】:
不幸的是,截至 2017 年 8 月,Pandas DataFame.apply() 仍仅限于使用单核,这意味着当您运行 df.apply(myfunc, axis=1) 时,多核机器将浪费其大部分计算时间。
如何使用所有内核在数据帧上并行运行应用程序?
【问题讨论】:
不幸的是,截至 2017 年 8 月,Pandas DataFame.apply() 仍仅限于使用单核,这意味着当您运行 df.apply(myfunc, axis=1) 时,多核机器将浪费其大部分计算时间。
如何使用所有内核在数据帧上并行运行应用程序?
【问题讨论】:
最简单的方法是使用Dask's map_partitions。你需要这些导入(你需要pip install dask):
import pandas as pd
import dask.dataframe as dd
from dask.multiprocessing import get
语法是
data = <your_pandas_dataframe>
ddata = dd.from_pandas(data, npartitions=30)
def myfunc(x,y,z, ...): return <whatever>
res = ddata.map_partitions(lambda df: df.apply((lambda row: myfunc(*row)), axis=1)).compute(get=get)
(如果你有 16 个核心,我认为 30 个分区是合适的)。为了完整起见,我在我的机器(16 核)上计时了差异:
data = pd.DataFrame()
data['col1'] = np.random.normal(size = 1500000)
data['col2'] = np.random.normal(size = 1500000)
ddata = dd.from_pandas(data, npartitions=30)
def myfunc(x,y): return y*(x**2+1)
def apply_myfunc_to_DF(df): return df.apply((lambda row: myfunc(*row)), axis=1)
def pandas_apply(): return apply_myfunc_to_DF(data)
def dask_apply(): return ddata.map_partitions(apply_myfunc_to_DF).compute(get=get)
def vectorized(): return myfunc(data['col1'], data['col2'] )
t_pds = timeit.Timer(lambda: pandas_apply())
print(t_pds.timeit(number=1))
28.16970546543598
t_dsk = timeit.Timer(lambda: dask_apply())
print(t_dsk.timeit(number=1))
2.708152851089835
t_vec = timeit.Timer(lambda: vectorized())
print(t_vec.timeit(number=1))
0.010668013244867325
从 pandas apply 到 dask apply 在分区上提供 10 倍加速。当然,如果你有一个可以向量化的函数,你应该——在这种情况下,函数 (y*(x**2+1)) 被简单地向量化了,但是有很多东西是不可能向量化的。
【讨论】:
The get= keyword has been deprecated. Please use the scheduler= keyword instead with the name of the desired scheduler like 'threads' or 'processes'
ValueError: cannot reindex from a duplicate axis。要解决这个问题,您应该通过df = df[~df.index.duplicated()] 删除重复索引或通过df.reset_index(inplace=True) 重置索引。
您可以使用swifter 包:
pip install swifter
(请注意,您可能希望在 virtualenv 中使用它以避免与已安装的依赖项发生版本冲突。)
Swifter 作为 pandas 的插件,让您可以重用 apply 函数:
import swifter
def some_function(data):
return data * 10
data['out'] = data['in'].swifter.apply(some_function)
它会自动找出并行化函数的最有效方法,无论它是否被矢量化(如上例所示)。
More examples 和 performance comparison 在 GitHub 上可用。请注意,该软件包正在积极开发中,因此 API 可能会发生变化。
另请注意,此will not work automatically 用于字符串列。当使用字符串时,Swifter 将回退到一个“简单”的 Pandas apply,它不会是并行的。在这种情况下,即使强制它使用 dask 也不会提高性能,最好手动拆分数据集并使用 parallelizing using multiprocessing。
【讨论】:
allow_dask_on_strings(enable=True):df.swifter.allow_dask_on_strings(enable=True).apply(some_function) 来源:github.com/jmcarpenter2/swifter/issues/45
你可以试试pandarallel:一个简单而高效的工具,可以在你的所有 CPU 上并行化你的 pandas 操作(在 Linux 和 macOS 上)
from pandarallel import pandarallel
from math import sin
pandarallel.initialize()
# FORBIDDEN
df.parallel_apply(lambda x: sin(x**2), axis=1)
# ALLOWED
def func(x):
return sin(x**2)
df.parallel_apply(func, axis=1)
【讨论】:
如果你想留在原生python:
import multiprocessing as mp
with mp.Pool(mp.cpu_count()) as pool:
df['newcol'] = pool.map(f, df['col'])
将以并行方式将函数 f 应用于数据帧 df 的列 col
【讨论】:
pandas/core/frame.py 中得到了来自__setitem__ 的ValueError: Length of values does not match length of index。不确定我是否做错了什么,或者分配给df['newcol'] 是否不是线程安全的。
只想给Dask一个更新答案
import dask.dataframe as dd
def your_func(row):
#do something
return row
ddf = dd.from_pandas(df, npartitions=30) # find your own number of partitions
ddf_update = ddf.apply(your_func, axis=1).compute()
在我的 100,000 条记录中,没有 Dask:
CPU时间:用户6分32秒,系统:100毫秒,总计:6分32秒 挂壁时间:6分32秒
与 Dask:
CPU 时间:用户 5.19 秒,系统:784 毫秒,总计:5.98 秒 挂墙时间:1分3秒
【讨论】:
要使用所有(物理或逻辑)内核,您可以尝试使用mapply 替代swifter 和pandarallel。
您可以在初始化时设置核心数量(和分块行为):
import pandas as pd
import mapply
mapply.init(n_workers=-1)
...
df.mapply(myfunc, axis=1)
默认情况下 (n_workers=-1),软件包使用系统上所有可用的物理 CPU。如果您的系统使用超线程(通常会显示两倍的物理 CPU 数量),mapply 将产生一个额外的工作人员来优先处理多处理池而不是系统上的其他进程。
根据您对all your cores 的定义,您也可以改用所有逻辑内核(请注意,像这样受 CPU 限制的进程将争夺物理 CPU,这可能会减慢您的操作速度):
import multiprocessing
n_workers = multiprocessing.cpu_count()
# or more explicit
import psutil
n_workers = psutil.cpu_count(logical=True)
【讨论】:
这是一个 sklearn 基础转换器的示例,其中 pandas 应用是并行化的
import multiprocessing as mp
from sklearn.base import TransformerMixin, BaseEstimator
class ParllelTransformer(BaseEstimator, TransformerMixin):
def __init__(self,
n_jobs=1):
"""
n_jobs - parallel jobs to run
"""
self.variety = variety
self.user_abbrevs = user_abbrevs
self.n_jobs = n_jobs
def fit(self, X, y=None):
return self
def transform(self, X, *_):
X_copy = X.copy()
cores = mp.cpu_count()
partitions = 1
if self.n_jobs <= -1:
partitions = cores
elif self.n_jobs <= 0:
partitions = 1
else:
partitions = min(self.n_jobs, cores)
if partitions == 1:
# transform sequentially
return X_copy.apply(self._transform_one)
# splitting data into batches
data_split = np.array_split(X_copy, partitions)
pool = mp.Pool(cores)
# Here reduce function - concationation of transformed batches
data = pd.concat(
pool.map(self._preprocess_part, data_split)
)
pool.close()
pool.join()
return data
def _transform_part(self, df_part):
return df_part.apply(self._transform_one)
def _transform_one(self, line):
# some kind of transformations here
return line
【讨论】:
self._preprocess_part?我只找到_transform_part
这里另一个使用 Joblib 和一些来自 scikit-learn 的帮助代码。轻量级(如果您已经拥有 scikit-learn),如果您希望更好地控制它正在做什么,那就太好了,因为 joblib 很容易被破解。
from joblib import parallel_backend, Parallel, delayed, effective_n_jobs
from sklearn.utils import gen_even_slices
from sklearn.utils.validation import _num_samples
def parallel_apply(df, func, n_jobs= -1, **kwargs):
""" Pandas apply in parallel using joblib.
Uses sklearn.utils to partition input evenly.
Args:
df: Pandas DataFrame, Series, or any other object that supports slicing and apply.
func: Callable to apply
n_jobs: Desired number of workers. Default value -1 means use all available cores.
**kwargs: Any additional parameters will be supplied to the apply function
Returns:
Same as for normal Pandas DataFrame.apply()
"""
if effective_n_jobs(n_jobs) == 1:
return df.apply(func, **kwargs)
else:
ret = Parallel(n_jobs=n_jobs)(
delayed(type(df).apply)(df[s], func, **kwargs)
for s in gen_even_slices(_num_samples(df), effective_n_jobs(n_jobs)))
return pd.concat(ret)
用法:result = parallel_apply(my_dataframe, my_func)
【讨论】:
由于问题是“如何使用所有内核在数据帧上并行运行应用程序?”,答案也可以是modin。您可以并行运行所有内核,但实时性更差。
见https://github.com/modin-project/modin。它运行在dask 或ray 的顶部。他们说“Modin 是为 1MB 到 1TB+ 的数据集设计的 DataFrame。”我试过了:pip3 install "modin"[ray]"。 Modin vs pandas 是 - 六核 12 秒 vs. 6 秒。
【讨论】:
而不是
df["new"] = df["old"].map(fun)
做
from joblib import Parallel, delayed
df["new"] = Parallel(n_jobs=-1, verbose=10)(delayed(fun)(i) for i in df["old"])
对我来说这是一个轻微的改进
import multiprocessing as mp
with mp.Pool(mp.cpu_count()) as pool:
df["new"] = pool.map(fun, df["old"])
如果作业非常小,您将获得进度指示和自动批处理。
【讨论】:
本机 Python 解决方案(使用 numpy),可以按照原始问题的要求应用于整个 DataFrame(不仅在单个列上)
import numpy as np
import multiprocessing as mp
dfs = np.array_split(df, 8000) # divide the dataframe as desired
def f_app(df):
return df.apply(myfunc, axis=1)
with mp.Pool(mp.cpu_count()) as pool:
res = pd.concat(pool.map(f_app, dfs))
【讨论】: