【问题标题】:Parellel function call in pythonpython中的并行函数调用
【发布时间】:2018-06-13 08:28:59
【问题描述】:

我对 python 很陌生。我一直在考虑将以下代码用于并行调用,其中在 lambda 的帮助下格式化 doj 值列表,

m_df[['doj']] = m_df[['doj']].apply(lambda x: formatdoj(*x), axis=1)

def formatdoj(doj):
    doj = str(doj).split(" ")[0]
    doj = datetime.strptime(doj, '%Y' + "-" + '%m' + "-" + "%d")
    return doj

由于列表有数百万条记录,因此格式化所有记录所需的时间很长。

如何在python中进行类似于c#中Parellel.Foreach的parellel函数调用?

【问题讨论】:

标签: python


【解决方案1】:

我认为在您的情况下使用并行计算有点过头了。缓慢来自代码,而不是使用单个处理器。我将通过一些步骤向您展示如何使其更快,猜测您正在使用 Pandas 数据框以及您的数据框包含什么(请遵守 SO 指南并包含一个完整的工作示例!!)

对于我的测试,我使用了以下具有 100k 行的随机数据帧(按比例放大以适应您的情况):

N=int(1e5)
m_df = pd.DataFrame([['{}-{}-{}'.format(y,m,d)]
                        for y,m,d in zip(np.random.randint(2007,2019,N),
                        np.random.randint(1,13,N),
                        np.random.randint(1,28,N))],
                    columns=['doj'])

现在这是你的代码:

tstart = time()
m_df[['doj']] = m_df[['doj']].apply(lambda x: formatdoj(*x), axis=1)
print("Done in {:.3f}s".format(time()-tstart))

在我的机器上运行大约需要 5.1 秒。它有几个问题。第一个是您使用的是数据框而不是系列,尽管您只处理一列,并创建了一个无用的 lambda 函数。简单地做:

m_df['doj'].apply(formatdoj)

将时间缩短至 1.6 秒。在 python 中,使用 '+' 连接字符串也很慢,您可以将 formatdoj 更改为:

def faster_formatdoj(doj):
    return datetime.strptime(doj.split()[0], '%Y-%m-%d')
m_df['doj'] = m_df['doj'].apply(faster_formatdoj)

这不是一个很大的改进,但确实将时间缩短到了 1.5 秒。如果你需要真正加入字符串(因为它们不是固定的),而是使用'-'.join('%Y','%m','%d'),这样会更快。

但真正的瓶颈来自于多次使用 datetime.strptime。它本质上是一个缓慢的命令 - 日期是一个笨重的东西。另一方面,如果你有数百万个日期,并假设它们自人类诞生以来就没有均匀分布,那么它们很可能被大量复制。因此,您应该真正做到以下几点:

tstart = time()
# Create a new column with only the first word
m_df['doj_split'] = m_df['doj'].apply(lambda x: x.split()[0])
converter = {
    x: faster_formatdoj(x) for x in m_df['doj_split'].unique()
}
m_df['doj'] = m_df['doj_split'].apply(lambda x: converter[x])
# Drop the column we added
m_df.drop(['doj_split'], axis=1, inplace=True)
print("Done in {:.3f}s".format(time()-tstart))

这大约在 0.2/0.3 秒内运行,比您的原始实现快 10 倍以上。

毕竟,如果您仍然运行缓慢,您可以考虑并行工作(而不是单独并行化第一个“拆分”指令,也许还有 apply-lambda 部分,否则您将创建许多不同的“转换器”字典使增益无效)。但我会把它作为最后一步,而不是第一个解决方案......

[编辑]:最初在最后一个代码框的第一步中,我使用了m_df['doj_split'] = m_df['doj'].str.split().apply(lambda x: x[0]),它在功能上是等效的,但比m_df['doj_split'] = m_df['doj'].apply(lambda x: x.split()[0]) 慢一点。我不完全确定为什么,可能是因为它本质上是应用了两个函数而不是一个。

【讨论】:

    【解决方案2】:

    最好的办法是使用dask。 Dask 有一个 data_frame 类型,您可以使用它来创建类似的数据帧,但是,在执行计算功能时,您可以使用 num_worker 参数指定核心数。这将使任务并行化

    【讨论】:

      【解决方案3】:

      由于我不确定您的示例,我将使用multiprocessing 库为您提供另一个示例:

      # -*- coding: utf-8 -*-
      import multiprocessing as mp
      
      input_list = ["str1", "str2", "str3", "str4"]
      
      def format_str(str_input):
          str_output = str_input + "_test"
          return str_output
      
      if __name__ == '__main__':
          with mp.Pool(processes = 2) as p:
              result = p.map(format_str, input_list)
      
          print (result)
      

      现在,假设您要映射一个带有多个参数的函数,那么您应该使用starmap()

      # -*- coding: utf-8 -*-
      import multiprocessing as mp
      
      input_list = ["str1", "str2", "str3", "str4"]
      
      def format_str(str_input, i):
          str_output = str_input + "_test" + str(i)
          return str_output
      
      if __name__ == '__main__':
          with mp.Pool(processes = 2) as p:
              result = p.starmap(format_str, [(input_list, i) for i in range(len(input_list))])
      
          print (result)
      

      不要忘记将 Pool 放在 if __name__ == '__main__': 中,并且 multiprocessing 将无法在诸如 spyder(或其他)之类的 IDE 中工作,因此您需要在 cmd 中运行脚本。

      要保留结果,您可以将它们保存到文件中,或者在末尾使用os.system("pause") (Windows) 或 Linux 上的 input() 保持 cmd 打开。

      这是一种在 python 中使用多处理的相当简单的方法。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2013-09-17
        • 1970-01-01
        • 1970-01-01
        • 2012-12-11
        • 1970-01-01
        • 1970-01-01
        • 2014-06-05
        相关资源
        最近更新 更多