【问题标题】:Pandas df.iterrows() parallelizationPandas df.iterrows() 并行化
【发布时间】:2017-03-14 10:28:41
【问题描述】:

我想并行化以下代码:

for row in df.iterrows():
    idx = row[0]
    k = row[1]['Chromosome']
    start,end = row[1]['Bin'].split('-')

    sequence = sequence_from_coordinates(k,1,start,end) #slow download form http

    df.set_value(idx,'GC%',gc_content(sequence,percent=False,verbose=False))
    df.set_value(idx,'G4 repeats', sum([len(list(i)) for i in g4_scanner(sequence)]))
    df.set_value(idx,'max flexibility',max([item[1] for item in dna_flex(sequence,verbose=False)]))

我尝试使用multiprocessing.Pool(),因为每一行都可以独立处理,但我不知道如何共享DataFrame。我也不确定这是与 pandas 进行并行化的最佳方法。有什么帮助吗?

【问题讨论】:

  • 默认情况下,您的逐行迭代很慢。您可以尝试找到一种方法来矢量化您的操作并且无需迭代即可完成,或者您将数据帧拆分为几个大块并并行迭代每个块。
  • 当然,这是一种方法。但我仍在寻找更好的方法,如果它存在的话。
  • 你考虑过使用 dask 吗?它会为您完成大部分并行化
  • 我不知道Dask,我去看看。

标签: python pandas multiprocessing


【解决方案1】:

正如@Khris 在他的评论中所说,你应该将你的数据框分成几个大块并并行迭代每个块。您可以将数据帧任意拆分为随机大小的块,但根据您计划使用的进程数将数据帧分成大小相等的块更有意义。幸运的是,其他人为我们提供了already figured out how to do that part

# don't forget to import
import pandas as pd
import multiprocessing

# create as many processes as there are CPUs on your machine
num_processes = multiprocessing.cpu_count()

# calculate the chunk size as an integer
chunk_size = int(df.shape[0]/num_processes)

# this solution was reworked from the above link.
# will work even if the length of the dataframe is not evenly divisible by num_processes
chunks = [df.iloc[df.index[i:i + chunk_size]] for i in range(0, df.shape[0], chunk_size)]

这将创建一个列表,其中包含我们的数据框。现在我们需要将它与一个操作数据的函数一起传递到我们的池中。

def func(d):
   # let's create a function that squares every value in the dataframe
   return d * d

# create our pool with `num_processes` processes
pool = multiprocessing.Pool(processes=num_processes)

# apply our function to each chunk in the list
result = pool.map(func, chunks)

此时,result 将是一个列表,在每个块被操作后保存它。在这种情况下,所有值都已平方。现在的问题是原始数据框尚未修改,因此我们必须将其所有现有值替换为我们池中的结果。

for i in range(len(result)):
   # since result[i] is just a dataframe
   # we can reassign the original dataframe based on the index of each chunk
   df.iloc[result[i].index] = result[i]

现在,我的数据帧操作函数已被矢量化,如果我只是将其应用于整个数据帧而不是拆分成块,可能会更快。但是,在您的情况下,您的函数将遍历每个块的每一行,然后返回该块。这允许您一次处理num_process 行。

def func(d):
   for row in d.iterrow():
      idx = row[0]
      k = row[1]['Chromosome']
      start,end = row[1]['Bin'].split('-')

      sequence = sequence_from_coordinates(k,1,start,end) #slow download form http
      d.set_value(idx,'GC%',gc_content(sequence,percent=False,verbose=False))
      d.set_value(idx,'G4 repeats', sum([len(list(i)) for i in g4_scanner(sequence)]))
      d.set_value(idx,'max flexibility',max([item[1] for item in dna_flex(sequence,verbose=False)]))
   # return the chunk!
   return d

然后您重新分配原始数据框中的值,并且您已成功并行化此过程。

我应该使用多少个进程?

您的最佳表现将取决于此问题的答案。而“所有的过程!!!!”是一个答案,一个更好的答案更加细微。在某一点之后,在一个问题上投入更多的进程实际上会产生比其价值更多的开销。这被称为Amdahl's Law。同样,我们很幸运,其他人已经为我们解决了这个问题:

  1. Python multiprocessing's Pool process limit
  2. How many processes should I run in parallel?

一个好的默认值是使用multiprocessing.cpu_count(),这是multiprocessing.Pool 的默认行为。 According to the documentation "如果 processes 为 None 则使用 cpu_count() 返回的数字。"这就是为什么我在开头将num_processes设置为multiprocessing.cpu_count()。这样,如果您迁移到更强大的机器,您就可以从中受益,而无需直接更改 num_processes 变量。

【讨论】:

  • 如果 pandas 显示警告,请使用 chunks = [df.iloc[i:i + chunk_size,:] for i in range(0, df.shape[0], chunk_size)]
  • np.array_split() 可能是比chunks = [df.ix[df.index[i:i + chunk_size]] for i in range(0, df.shape[0], chunk_size)] 更好的选择。前者会自动处理行数不能整除的情况,语法稍微简单一些。
【解决方案2】:

一种更快的方法(在我的情况下约为 10%):

与已接受答案的主要区别: 使用pd.concatnp.array_split 拆分和加入数据帧。

import multiprocessing
import numpy as np


def parallelize_dataframe(df, func):
    num_cores = multiprocessing.cpu_count()-1  #leave one free to not freeze machine
    num_partitions = num_cores #number of partitions to split dataframe
    df_split = np.array_split(df, num_partitions)
    pool = multiprocessing.Pool(num_cores)
    df = pd.concat(pool.map(func, df_split))
    pool.close()
    pool.join()
    return df

其中func 是您要应用于df 的函数。使用partial(func, arg=arg_val) 获得不止一个论点。

【讨论】:

  • 只是好奇,pool.map 是否维护数据帧的顺序。换句话说,pool.map 的输出是否与传入的块的顺序相同?如果不是,那么pd.concat 可能不会以原始顺序重建数据框。我不知道np.aray_split,但我并不惊讶它更快。 pd.concat 也可能比使用 df.ix 重新分配更快
  • @Jalepeno112 是的,据我所知,数据框以正确的顺序重新组合在一起。我不知道是否有办法执行它,但我有时间序列数据,它还没有引起问题。尽管由于我的索引是时间戳,但如果订单混乱,再次对它们进行排序应该不是问题。我发现的另一个技巧是使用 itertuples(),它又快了 30%。
  • 这拯救了我的一天。非常感谢@ic_fl2
  • 你能帮我回答这个问题吗:- stackoverflow.com/questions/53561794/…
  • 这是一个非常好的答案!
【解决方案3】:

考虑使用 dask.dataframe,例如此示例中显示了类似问题:https://stackoverflow.com/a/53923034/4340584

import dask.dataframe as ddf
df_dask = ddf.from_pandas(df, npartitions=4)   # where the number of partitions is the number of cores you want to use
df_dask['output'] = df_dask.apply(lambda x: your_function(x), meta=('str')).compute(scheduler='multiprocessing')

【讨论】:

  • dask 解决方案看起来比在pandas 中手动并行化计算要简单得多!
【解决方案4】:

要在数据帧的分区上使用 Dask(而不是在 axis 上运行的 dask.apply),您可以使用 map_partitions

import multiprocessing
import dask.dataframe as ddf

# get num cpu cores
num_partitions = multiprocessing.cpu_count()

# create dask DF
df_dask = ddf.from_pandas(your_dataframe, npartitions=num_partitions)

# apply func to every partition in parallel
output = df_dask.map_partitions(func, meta=('output_col1_type','output_col2_type')).compute(scheduler='multiprocessing')

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-01-10
    • 2020-04-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-01-01
    • 2014-11-29
    • 1970-01-01
    相关资源
    最近更新 更多