存储结果ThreadPoolExecutor

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)

aba*_*ert 5

当您调用时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)