正如@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。同样,我们很幸运,其他人已经为我们解决了这个问题:
- Python multiprocessing's Pool process limit
- 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 变量。