【问题标题】:pandas multiprocessing apply熊猫多处理应用
【发布时间】: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


【解决方案1】:

您可以使用https://github.com/nalepae/pandarallel,如下例所示:

from pandarallel import pandarallel
from math import sin

pandarallel.initialize()

def func(x):
    return sin(x**2)

df.parallel_apply(func, axis=1)

【讨论】:

  • 这个答案应该会得到更多的支持。速度非常快。
  • 此解决方案原生适用于 linux 和 macOS。在 Windows 上,Pandaral·lel 仅在 Python 会话从 Windows Subsystem for Linux (WSL) 执行时才有效。
  • 在 Windows 上,我收到此错误:ValueError: cannot find context for 'fork'
【解决方案2】:

基于作者解决方案的更通用版本,允许在每个函数和数据帧上运行它:

from multiprocessing import  Pool
from functools import partial
import numpy as np

def parallelize(data, func, num_of_processes=8):
    data_split = np.array_split(data, num_of_processes)
    pool = Pool(num_of_processes)
    data = pd.concat(pool.map(func, data_split))
    pool.close()
    pool.join()
    return data

def run_on_subset(func, data_subset):
    return data_subset.apply(func, axis=1)

def parallelize_on_rows(data, func, num_of_processes=8):
    return parallelize(data, partial(run_on_subset, func), num_of_processes)

所以下面一行:

df.apply(some_func, axis=1)

会变成:

parallelize_on_rows(df, some_func) 

【讨论】:

  • 带参数的some_func怎么样?
  • @AlaaM。 - 你可以使用 partial 。 docs.python.org/2/library/functools.html#functools.partial
  • @TomRaz 在这种情况下,当我通常会做这样的事情时,我该如何使用部分? dataframe.apply(lambda row: process(row.attr1, row.attr2, ...))
  • @frei - lambda 函数不能用于多处理,因为它们不能被腌制。在此处查看更多信息:stackoverflow.com/a/8805244/1781490您可以改用普通功能吗?
  • 我明白了。这就是我需要知道我是否应该重构整个方法的部分
【解决方案3】:

这是我发现有用的一些代码。自动将数据帧拆分为您拥有的多个 cpu 核心。

import pandas as pd
import numpy as np
import multiprocessing as mp

def parallelize_dataframe(df, func):
    num_processes = mp.cpu_count()
    df_split = np.array_split(df, num_processes)
    with mp.Pool(num_processes) as p:
        df = pd.concat(p.map(func, df_split))
    return df

def parallelize_function(df):
    df[column_output] = df[column_input].apply(example_function)
    return df

def example_function(x):
    x = x*2
    return x

运行:

df_output = parallelize_dataframe(df, parallelize_function)

【讨论】:

    【解决方案4】:

    由于我没有太多你的数据脚本,这是一个猜测,但我建议在回调中使用 p.map 而不是 apply_async

    p = mp.Pool(8)
    pool_results = p.map(process, np.array_split(big_df,8))
    p.close()
    p.join()
    results = []
    for result in pool_results:
        results.extend(result)
    

    【讨论】:

    • 如果 name == 'main',我必须将调用放入其中。通过其他小的更改,我设法使其工作,但是我不确定池结果中的结果数据帧是否以与拆分相同的顺序返回。我必须检查一下。
    • 看这里daskstackoverflow.com/questions/37979167/…的解决方案
    【解决方案5】:

    这对我很有效:

    rows_iter = (row for _, row in df.iterrows())
    
    with multiprocessing.Pool() as pool:
        df['new_column'] = pool.map(process_apply, rows_iter)
    

    【讨论】:

      【解决方案6】:

      要使用所有(物理或逻辑)内核,您可以尝试使用mapply 作为swifterpandarallel 的替代方案。

      您可以在初始化时设置核心数量(和分块行为):

      import pandas as pd
      import mapply
      
      mapply.init(n_workers=-1)
      
      def process_apply(x):
          # do some stuff to data here
      
      def process(df):
          # spawns a pathos.multiprocessing.ProcessPool if sensible
          res = df.mapply(process_apply, axis=1)
          return res
      

      默认情况下 (n_workers=-1),软件包使用系统上所有可用的物理 CPU。如果您的系统使用超线程(通常会显示两倍的物理 CPU 数量),mapply 将产生一个额外的工作人员来优先处理多处理池而不是系统上的其他进程。

      您也可以改为使用所有逻辑内核(请注意,像这样受 CPU 限制的进程将争夺物理 CPU,这可能会减慢您的操作速度):

      import multiprocessing
      n_workers = multiprocessing.cpu_count()
      
      # or more explicit
      import psutil
      n_workers = psutil.cpu_count(logical=True)
      

      【讨论】:

        【解决方案7】:

        当我使用multiprocessing.map() 将函数应用于大型数据帧的不同块时,我也遇到了同样的问题。

        我只想补充几点,以防其他人遇到和我一样的问题。

        1. 记得加if __name__ == '__main__':
        2. .py文件中执行文件,如果使用ipython/jupyter notebook,则无法运行multiprocessing(我的情况确实如此,虽然我不知道)

        【讨论】:

          【解决方案8】:

          安装Pyxtension,它简化了并行映射的使用并像这样使用:

          from pyxtension.streams import stream
          
          big_df = pd.concat(stream(np.array_split(df, multiprocessing.cpu_count())).mpmap(process))
          

          【讨论】:

            【解决方案9】:

            我最终使用 concurrent.futures.ProcessPoolExecutor.map 代替了 multiprocessing.Pool.map,对于一些需要 12 秒的串行代码,这需要 316 微秒。

            【讨论】:

              【解决方案10】:

              Python 的pool.starmap() method 也可用于简洁地将并行性引入apply 列值作为参数传递的用例,例如:

              df.apply(lambda row: my_func(row["col_1"], row["col_2"], ...), axis=1)
              

              完整示例和基准测试:

              import time
              from multiprocessing import Pool
              
              import numpy as np
              import pandas as pd
              
              
              def mul(a, b, c):
                  # For illustration, could obviously be vectorized
                  return a * b * c
              
              df = pd.DataFrame(np.random.randint(0, 100, size=(10_000_000, 3)), columns=list('ABC'))
              
              # Standard apply
              start = time.time()
              df["mul"] = df.apply(lambda row: mul(row["A"], row["B"], row["C"]), axis=1)
              print(f"Standard apply took {time.time() - start:.0f} seconds.") 
              
              # Starmap apply
              start = time.time()
              with Pool(10) as pool:
                  df["mul_pool"] = pool.starmap(mul, zip(df["A"], df["B"], df["C"]))
              print(f"Starmap apply took {time.time() - start:.0f} seconds.") 
              
              pd.testing.assert_series_equal(df["mul"], df["mul_pool"], check_names=False)
              
              
              >>> Standard apply took 72 seconds.
              >>> Starmap apply took 5 seconds.
              

              这样做的好处是不依赖外部库,而且可读性强。

              【讨论】:

                猜你喜欢
                • 1970-01-01
                • 1970-01-01
                • 2022-06-12
                • 2016-02-06
                • 2018-10-20
                • 1970-01-01
                • 1970-01-01
                相关资源
                最近更新 更多