【问题标题】:Multiprocessing vs Concurrent.futures library python (not working on Google Compute Engine)Multiprocessing vs Concurrent.futures library python(不适用于Google Compute Engine)
【发布时间】:2019-07-11 12:40:27
【问题描述】:

我正在尝试并行化一个 pandas 操作,它将具有逗号分隔值的数据框列拆分为 2 列。在我的 python 实例上,正常的 pandas 操作大约需要 5 秒,它直接在该特定列上使用 df.str.split。我的数据框包含 200 万行,因此我试图降低代码运行时间。

作为第一种并行化方法,我正在使用 Python 的多处理库,方法是创建与我的实例上可用的 CPU 内核数量相等的池。对于同一问题的第二种方法,我使用 concurrent.futures 库,提到了 4 的 chunksize。 但是,我看到多处理库与正常的 pandas 操作(5 秒)所用的时间大致相同,而 concurrent.futures 运行同一行的时间超过一分钟。

1) Google Compute Engine 是否支持这些 Python 多处理库? 2) 为什么并行处理在 GCP 上不起作用?

提前致谢。下面是示例代码:

import pandas as pd
from multiprocessing import Pool

def split(e):
    return e.split(",")

df =  pd.DataFrame({'XYZ':['CAT,DOG', 
      'CAT,DOG','CAT,DOG']})

pool = Pool(4)
df_new = pd.DataFrame(pool.map(split, df['XYZ'], columns = ['a','b'])
df_new = pd.concat([df, df_new], axis=1)

上面的代码与下面的代码所用的时间大致相同,下面的代码是只使用一个核心的普通 pandas 操作:

df['a'], df['b'] = df['XYZ'].str.split(',',1).str

使用 concurrent.futures:

import concurrent.futures
with concurrent.futures.ProcessPoolExecutor() as pool:
     a = pd.DataFrame(pool.map(split, df['XYZ'], chunksize = 4), 
     columns=['a','b'])
print (a)

上面使用 concurrent.futures 的代码在 GCP 上运行需要一分钟多的时间。请注意,我发布的代码只是示例代码。我在项目中使用的数据框有 200 万行这样的行。任何帮助将不胜感激!

【问题讨论】:

  • 嗨,我们说的是哪个产品呢?请记住,GCP 是一个包含大量产品(GCE、GKE DataProc)的平台。你能说得更具体些吗?
  • 它是一个 4 cpus 的 GCE。谢谢

标签: python-3.x google-cloud-platform multiprocessing google-compute-engine concurrent.futures


【解决方案1】:

为什么选择chunksize=4?这非常小,对于 200 万行,这只会将其分解为 500,000 次操作。总运行时间可能只需要 1/4 的时间,但额外的开销可能会比单线程方法花费更长的时间。

我建议使用更大的chunksize。 10,000 到 200,000 的任何值都可能合适,但您应该根据获得的结果进行一些实验来调整它。

【讨论】:

  • 嗨达斯汀,谢谢你的建议。当我将块大小增加到 200,000 左右时,结果会更好。然而,在用不同的值进行实验之后,我能得到的最佳时间是 6 秒左右。我认为这是因为并行化中涉及的子进程的开销。有什么解决方法吗?
  • 对我来说,6 秒对于 2M 行来说听起来很合理。我认为这里的操作很简单,使其并发可能不会有任何改进。
猜你喜欢
  • 2018-02-14
  • 2020-09-03
  • 2014-09-13
  • 2018-12-02
  • 2021-04-20
  • 1970-01-01
  • 1970-01-01
  • 2015-01-25
  • 1970-01-01
相关资源
最近更新 更多