【问题标题】:Parallelising for-loop using multiprocessing pool function使用多处理池函数并行化 for 循环
【发布时间】:2018-06-05 15:15:33
【问题描述】:

我试图按照这个位置的示例进行操作:

[How to use threading in Python?

我有一个这样的示例数据框 (df):

segment x_coord y_coord
a   1   1
a   2   4
a   1   7
b   2   3
b   4   3
b   8   3
c   4   4
c   2   5
c   7   8

并使用 for 循环为循环中的每个段创建 kd-tree,如下所示:

dist_name=df['segment'].unique()
for i in range(len(dist_name)):
    a=df[df['segment']==dist_name[i]]
    tree[i] = spatial.cKDTree(a[['x_coord','y_coord']])

如何使用以下链接中的示例并行创建树:

results = [] 
for url in urls:
  result = urllib2.urlopen(url)
  results.append(result)

并行化到 >>

pool = ThreadPool(4) 
results = pool.map(urllib2.urlopen, urls)

我的尝试

import pandas as pd
import time
from scipy import spatial
import random
from multiprocessing.dummy import Pool as ThreadPool 


dist_name=['a','b','c','d','e','f','g','h']

df=pd.DataFrame()

for i in range(len(dist_name)):
    if i==0:
       df['x_coord']=random.sample(range(1, 10000), 1000)
       df['y_coord']=random.sample(range(1, 10000), 1000)
       df['segment']=dist_name[i]
    else:
       tmp=pd.DataFrame()
       tmp['x_coord']=random.sample(range(1, 10000), 1000)
       tmp['y_coord']=random.sample(range(1, 10000), 1000)
       tmp['segment']=dist_name[i]
       df=df.append(tmp)



start_time = time.time()
for i in range(len(dist_name)):
    a=df[df['segment']==dist_name[i]]
    tree = spatial.cKDTree(a[['x_coord','y_coord']])

print("--- %s seconds ---" % (time.time() - start_time))

--- 0.0312347412109375 秒 ---

def func(name):
    a = df[df['segment'] == name]
    return spatial.cKDTree(a[['x_coord','y_coord']])

pool = ThreadPool(4) 

start_time = time.time()
tree = pool.map(func, dist_name)
print("--- %s seconds ---" % (time.time() - start_time))

--- 0.031250953674316406 秒 ---

【问题讨论】:

  • 您能否通过尝试将链接中的答案应用于您的代码来更新您的问题?您需要考虑要编写什么函数(例如func()),以便您可以编写pool.map(func, dist_name)
  • 用我的方法更新了

标签: python for-loop parallel-processing kdtree


【解决方案1】:

您的代码:

dist_name=df['segment'].unique()
for i in range(len(dist_name)):
    a=df[df['segment']==dist_name[i]]
    tree[i] = spatial.cKDTree(a[['x_coord','y_coord']])

需要转化为:

dist_name=df['segment'].unique()

def func(name):
    a = df[df['segment'] == name]
    return spatial.cKDTree(a[['x_coord','y_coord']])

还有你给pool.map的电话:

pool = ThreadPool(4) 
tree = pool.map(func, dist_name)

【讨论】:

  • 与正常运行相比,我没有观察到使用线程的任何性能改进。这里的罪魁祸首可能是什么?
  • 我更新了答案以表明没有性能差异。不确定,是什么导致了这种行为
  • 你为什么使用multiprocessing.dummy
  • 这是根据链接
  • 该链接还说这使用了幕后的线程库。这不会加快 CPU 绑定任务的速度。删除dummy,您应该进行多处理。请注意,对于一个只需要 31 毫秒的任务来说,这可能看起来更慢。您还需要if __name__ == '__main__': 成语。
猜你喜欢
  • 2022-01-12
  • 1970-01-01
  • 1970-01-01
  • 2017-01-28
  • 1970-01-01
  • 2022-01-03
  • 1970-01-01
  • 1970-01-01
  • 2021-08-15
相关资源
最近更新 更多