【问题标题】:Correct way of creating and managing dependencies between concurrent futures in Python在 Python 中创建和管理并发期货之间依赖关系的正确方法
【发布时间】:2021-12-24 15:48:52
【问题描述】:

现在我正在探索并行性(多处理和多线程)。 Futures 似乎很受欢迎,所以我正在尝试找出是否可以在 Futures 之间创建依赖关系并在我自己的应用程序中使用它们。

我将使用我编写的以下代码作为示例:

from concurrent.futures import ThreadPoolExecutor
import os

def dependent_future():
    def thread_fn(x, y):
        res = x + y
        return res

    def thread_fn_with_dep(xy_future, power):
        xy = xy_future.result()
        res = pow(xy, power)
        return res

    with ThreadPoolExecutor(max_workers=os.cpu_count()) as executor:
        xy_future = executor.submit(thread_fn, 13, 34)
        xy_power_future = executor.submit(thread_fn_with_dep, xy_future, 2)


dependent_future()

我想知道这是否是让一个未来依赖于另一个未来的正确方法。据我了解,调用result() 会阻止执行(这里希望只是执行thread_fn_with_dep(...) 的线程),所以只有在thread_fn 完成之后才会执行thread_fn_with_dep。但是,我提出的以下其他示例(表示多项式的求和表达式)无法按预期工作:

def multidep_future():
    def polynom(fs):
        sum = 0
        for f in as_completed(fs):
            res, idx = f.result()
            print("[future_idx : " + str(idx) + "] : res = " + str(res) + ", sum_intermediate = " + str(sum) + "]")
            sum = sum + res

        return sum

    def mult(a_k, x, k):
        x_pow_res = pow(x, k)
        a_x_mult = a_k*x_pow_res
        print("a_" + str(k) + "*x^" + str(k) + " = " + str(a_k) + "*" + str(x) + "^" + str(k)
              + " = " + str(a_k) + "*" + str(x_pow_res)
              + " = " + str(a_x_mult))

        return pow(a_k*x, k), k

    with ThreadPoolExecutor(max_workers=os.cpu_count()) as executor:
        x = 1
        futures = [
            executor.submit(mult, k, x, k)
            for k in range(5)
        ]

        poly = executor.submit(polynom, futures)

        print(poly.result())


multidep_future()

输出是

a_0*x^0 = 0*1^0 = 0*1 = 0
a_1*x^1 = 1*1^1 = 1*1 = 1
a_2*x^2 = 2*1^2 = 2*1 = 2
a_3*x^3 = 3*1^3 = 3*1 = 3
a_4*x^4 = 4*1^4 = 4*1 = 4
[future_idx : 2] : res = 4, sum_intermediate = 0]
[future_idx : 3] : res = 27, sum_intermediate = 4]
[future_idx : 0] : res = 1, sum_intermediate = 31]
[future_idx : 1] : res = 1, sum_intermediate = 32]
[future_idx : 4] : res = 256, sum_intermediate = 33]
289

我期待的是

a_0*x^0 = 0*1^0 = 0*1 = 0
a_1*x^1 = 1*1^1 = 1*1 = 1
a_2*x^2 = 2*1^2 = 2*1 = 2
a_3*x^3 = 3*1^3 = 3*1 = 3
a_4*x^4 = 4*1^4 = 4*1 = 4
[future_idx : 2] : res = 2, sum_intermediate = 3]
[future_idx : 3] : res = 3, sum_intermediate = 6]
[future_idx : 0] : res = 0, sum_intermediate = 0]
[future_idx : 1] : res = 1, sum_intermediate = 1]
[future_idx : 4] : res = 4, sum_intermediate = 10]
10

因为(排序)

  a_0*x^0 + a_1*x^1 + a_2*x^2 + a_3*x^3 + a_4*x^4
= 0*x^0 + 1*x^1 + 2*x^2 + 3*x^3 + 4*x^4
= 0*1^0 + 1*1^1 + 2*1^2 + 3*1^3 + 4*1^4
= 0*0 + 1*1 + 2*1 + 3*1 + 4*1
= 0 + 1 + 2 + 3 + 4
= 1 + 2 + 3 + 4
= 3 + 3 + 4
= 6 + 4
= 10

我显然在这里遗漏了一些东西。

此外,我想知道如何管理多个依赖项,特别是如果存在更复杂的任务链,其中 - 由于一个任务失败 - 只能执行剩余任务的子集。这更像是一个“如果依赖任务卡住或失败,依赖任务会发生什么”的问题。

一个更高级的例子是图像处理。通常,初始阶段是将源图像转换为具有减少色彩空间的图像(例如灰度又名灰色单色)。之后,我们可以应用高斯去噪,然后我们可以将结果传递给各种程序,例如开/关、边缘检测、角点检测、特征检测、特征匹配、将一些结果写入文件、发送图像网络上的数据等。这引入了多个瓶颈,其中几个步骤本质上是并行的,但依赖于相同的输入。在这种情况下,我想使用期货,因为我编写的示例代码就是事情的完成方式。

【问题讨论】:

    标签: python concurrency dependencies future


    【解决方案1】:

    好吧,它会起作用的。这种方法的一个缺点是您限制了并行化,因为工作线程可能会等待未来,而另一个任务将在队列中可用。

    更好的方法是仅在前一个未来准备就绪时才安排下一个任务。您可以为此使用future.add_done_callback。但是,为复杂的工作流程编写此程序很乏味。

    总的来说,您首先必须问自己是否值得。函数之间的复杂依赖网络很难维护和调试。由于更好的缓存使用率,计算通常会通过 fork-join(例如,单个图像的并行处理)获得更好的加速。如果并行化过度,I/O 甚至可能会降级。

    总而言之,更好地掌握多线程并仅将期货和线程池用作一种工具。例如,我发现使用一个线程池进行 I/O 和一个线程池进行计算(如果有的话)要好得多。

    【讨论】:

    • 查看更新的问题。我也会做几个线程。但现在我正在寻找不同的方法来并行做事。例如,我有 3 个进程 - 一个生成另外两个进程,一个运行 GUI,另一个执行图像处理。图像处理进程包含一个线程来轮询由 GUI 进程加载的传入图像。我想将与不同处理任务相关的负载分配到多个线程中,然后每个线程将其结果放入multiprocessing.Queue 以供 GUI 可视化。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2011-06-21
    • 2020-07-20
    • 1970-01-01
    • 2017-01-15
    • 2014-07-17
    • 1970-01-01
    • 2014-10-15
    相关资源
    最近更新 更多