【问题标题】:python executor spawn tasks from done callback (recursively submit tasks)python 执行器从完成回调生成任务(递归提交任务)
【发布时间】:2019-01-23 12:47:45
【问题描述】:

我正在尝试根据已完成任务的结果提交进一步的任务:

with concurrent.futures.ThreadPoolExecutor() as executor:
    future = executor.submit(my_task)
    def callback(future):
        for another_task in future.result():
            future = executor.submit(another_task)
            future.add_done_callback(callback)
    future.add_done_callback(callback)

但我得到了:

RuntimeError: 关机后无法安排新的未来

让执行程序等待回调的最佳方法是什么?信号量?

如果将ThreadPoolExecutor 替换为ProcessPoolExecutor,理想情况下解决方案应该可以很好地转移。

【问题讨论】:

    标签: python multithreading multiprocessing


    【解决方案1】:

    callback 在单独的线程中执行。因此,当你的回调逻辑开始执行时,主循环很有可能会让上下文管理器关闭Executor。这就是你观察RuntimeError的原因。

    最简单的解决方法是按顺序运行您的逻辑。

    futures = []
    
    with concurrent.futures.ThreadPoolExecutor() as pool:
        futures.append(pool.submit(task))
    
        while futures:
            for future in concurrent.futures.as_completed(futures):
                futures.remove(future)
    
                for new_task in future.result():
                    futures.append(pool.submit(new_task))
    

    请注意,此代码可能会导致以指数方式向Executor 提交任务。

    【讨论】:

    • concurrent.futures.as_completed 会阻止吗?这会以某种方式导致延迟提交新任务吗?
    • 它会阻塞,直到您的 future 对象之一准备好。无论如何,您必须要做的是,因为要提交的作业由future 本身的结果提供。
    • 考虑我们从 1 个简单任务和 2 个两个困难任务开始,最简单的任务首先发布并分成 50 个新的简单任务,但在完成剩余任务之前我们不会到达它们2 项艰巨的任务(无论我们有多少线程)。
    • 我的意思是——我们可能会接触到它们,但不会接触到它们产生的那些。本质上,我们正在做一些有点像任务图上的 BFS 的事情,而不一定要限制线程/进程。
    • 我不确定您要达到的目标。不过,上面的逻辑非常简单。 concurrent.futures.as_completed 阻塞,直到给定的futures 之一准备好。无论哪个顺序都无所谓。因此上面的逻辑是正确的:它一直等到有更多的事情要做。当新任务到来时(作为完成未来的结果),它会将它们异步调度到pool。您不能这样做,因为您的问题具有顺序性:新任务依赖于旧任务来完成。
    【解决方案2】:

    此解决方案保证在任何给定点最大化处理。 sleep 的使用不是那么优雅,但这是我迄今为止最好的。

    with concurrent.futures.ThreadPoolExecutor() as executor:
        pending_tasks = 1
        future = executor.submit(my_task)
        def callback(future):
            nonlocal pending_tasks
            for another_task in future.result():
                pending_tasks += 1
                future = executor.submit(another_task)
                future.add_done_callback(callback)
            pending_tasks -= 1
        future.add_done_callback(callback)
        while pending_tasks:
            time.sleep(10)
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-05-13
      • 1970-01-01
      相关资源
      最近更新 更多