Mik*_*ike 1 python python-multithreading concurrent.futures
我对使用“ concurrent.futures”进行并行处理还很陌生,并且正在测试一些简单的实验。我编写的代码似乎有效,但是我不确定如何存储结果。我试图创建一个列表(“ futures”)并将结果附加到该列表中,但这会大大减慢该过程。我想知道是否有更好的方法可以做到这一点。谢谢。
import concurrent.futures
import time
couple_ods= []
futures=[]
dtab={}
for i in range(100):
for j in range(100):
dtab[i,j]=i+j/2
couple_ods.append((i,j))
avg_speed=100
def task(i):
origin=i[0]
destination=i[1]
time.sleep(0.01)
distance=dtab[origin,destination]/avg_speed
return distance
start1=time.time()
def main():
with concurrent.futures.ThreadPoolExecutor() as executor:
for number in couple_ods:
future=executor.submit(task,number)
futures.append(future.result())
if __name__ == '__main__':
main()
end1=time.time()
Run Code Online (Sandbox Code Playgroud)
当您调用时future.result(),它将阻塞直到值准备就绪。因此,在这里并行性并没有带来任何好处-您开始一个任务,等待它完成,开始另一个任务,等待它完成,等等。
当然,您的示例首先不会从线程中受益。您的任务除了CPU受限的Python计算外什么也不做,这意味着(至少在CPython,MicroPython和PyPy中,它们是唯一附带的完整实现concurrent.futures),GIL(全局解释器锁定)将防止以下情况中的一种您的线程无法一次执行。
希望您的真实程序有所不同。如果它正在做I / O绑定的工作(发出网络请求,读取文件等),或者使用诸如NumPy的扩展库来释放GIL来处理繁重的CPU,那么它将可以正常工作。但是否则,您将需要在ProcessPoolExecutor这里使用。
无论如何,您想要做的就是将future自身添加到列表中,因此在等待任何期货之前,您会获得所有期货的列表:
for number in couple_ods:
future=executor.submit(task,number)
futures.append(future)
Run Code Online (Sandbox Code Playgroud)
然后,在开始所有作业之后,您可以开始等待它们。有三个简单的选项,而一个复杂的选项则需要更多控制。
(1)您可以直接将它们循环播放,以按提交的顺序等待它们:
for future in futures:
result = future.result()
dostuff(result)
Run Code Online (Sandbox Code Playgroud)
(2)如果需要在完成任何工作之前等待它们完成,则可以致电wait:
futures, _ = concurrent.futures.wait(futures)
for future in futures:
result = future.result()
dostuff(result)
Run Code Online (Sandbox Code Playgroud)
(3)如果您希望每一个准备就绪就立即处理,即使它们出现故障,请使用as_completed:
for result in concurrent.futures.as_completed(futures):
dostuff(result)
Run Code Online (Sandbox Code Playgroud)
请注意,在文档中使用此功能的示例提供了一些方法来标识完成的任务。如果你需要,它可以简单的通过每一个索引,然后return index, real_result,然后就可以for index, result in …进行循环。
(4)如果您需要更多控制权,则可以遍历wait到目前为止所做的任何事情:
while futures:
done, futures = concurrent.futures.wait(concurrent.futures.FIRST_COMPLETED)
for future in done:
result = future.result()
dostuff(result)
Run Code Online (Sandbox Code Playgroud)
该示例的功能与相同as_completed,但是您可以在其上写一些较小的变体以执行不同的操作,例如等待所有操作完成,但如果出现异常则提早取消。
对于许多简单情况,您可以只使用map执行程序的方法来简化第一个选项。就像内置map函数一样,它对参数中的每个值调用一次函数,然后为您提供一些内容,您可以循环执行以相同顺序获得结果,但它是并行执行的。所以:
for result in executor.map(task, couple_ods):
dostuff(result)
Run Code Online (Sandbox Code Playgroud)