当轴 = 0 时,熊猫并行应用

Mis*_*lav 7 python parallel-processing pandas python-multiprocessing

我想在所有 Pandas 列上并行应用一些函数。例如,我想并行执行此操作:

def my_sum(x, a):
    return x + a


df = pd.DataFrame({'num_legs': [2, 4, 8, 0],
                   'num_wings': [2, 0, 0, 0]})
df.apply(lambda x: my_sum(x, 2), axis=0)
Run Code Online (Sandbox Code Playgroud)

我知道有一个swifter包,但它不支持axis=0应用:

NotImplementedError:Swifter 无法在大型数据集上执行 axis=0 应用。Dask 目前没有实现 axis=0 应用。更多详情请访问https://github.com/jmcarpenter2/swifter/issues/10

Dask 也不支持此功能axis=0(根据swifter 中的文档)。

我用谷歌搜索了几个来源,但找不到简单的解决方案。

不敢相信这在熊猫中如此复杂。

Arc*_*ast 0

您可以使用 dask 延迟接口来设置自定义工作流程:

import pandas as pd
import dask
import distributed

# start local cluster, by default one worker per core
client = distributed.Client() 

@dask.delayed
def my_sum(x, a):
    return x + a

df = pd.DataFrame({'num_legs': [2, 4, 8, 0],
                   'num_wings': [2, 0, 0, 0]})    

# Here, we mimic the apply command. However, we do not
# actually run any computation. Instead, that line of code 
# results in a list of delayed objects, which contain the 
# information what computation should be performed eventually
delayeds = [my_sum(df[column], 2) for column in df.columns]

# send the list of delayed objects to the cluster, which will 
# start computing the result in parallel. 
# It returns future objects, pointing to the computation while
# it is still running
futures = client.compute(delayeds)

# get all the results, as soon as they are ready. This returns 
# a list of pandas Series objects, each is one column of the 
# output dataframe
computed_columns = client.gather(futures)

# create dataframe out of individual columns
computed_df = pd.concat(computed_columns, axis = 1)
Run Code Online (Sandbox Code Playgroud)

或者,您也可以使用 dask 的多处理后端:

import pandas as pd
import dask

@dask.delayed
def my_sum(x, a):
    return x + a

df = pd.DataFrame({'num_legs': [2, 4, 8, 0],
                   'num_wings': [2, 0, 0, 0]})    

# same as above
delayeds = [my_sum(df[column], 2) for column in df.columns]

# run the computation using the dask's multiprocessing backend
computed_columns = dask.compute(delayeds, scheduler = 'processes')

# create dataframe out of individual columns
computed_df = pd.concat(computed_columns, axis = 1)
Run Code Online (Sandbox Code Playgroud)