【问题标题】:Graceful Termination of Worker Pool工作池的优雅终止
【发布时间】:2015-12-28 00:40:15
【问题描述】:

我想生成 X 数量的 Pool 工作人员,并给他们每个人 X% 的工作要做。我的问题是,这项工作大约需要 20 分钟才能耗尽,每个额外的进程运行时间更长,由于正在完成的计算类型,我的答案可能会在几分钟或几小时内找到。我想做的是为单个工作人员实现某种方式来“嘿,我找到了”并使用该信号来杀死池的其余部分并继续我的计算。

关键点:

  • 我尝试过回调,它们似乎在整个池完成之前不会在starmap_async 上运行。
  • 我只关心找到的第一个合适的答案。
  • 我不共享资源,并且意外进程死亡,尽管很粗鲁,但完全可以接受。

我也考虑过使用队列,但它不会,因为我传递给每个队列的工作范围已经内置到函数的参数中。

下面是我正在使用的一个非常枯燥的版本(我正在使用的计算可能需要数小时才能完成超过 42 亿个复杂的迭代。)

def doWork():
    workers = Pool(2)
    results = workers.starmap_async( func = distSearch , iterable = Sections1_5,  callback = killPool )
    workers.close()
    print("Found answer : {}".format(results.get()))
    workers.join()

def killPool():
    workers.terminate()
    print("Worker Pool Terminated")

我可能应该指定我的进程只有在找到答案时才返回,否则一旦完成就退出。我查看了this 线程,但我完全迷路了,而且在工作池的返回/回调中持续检查获胜条件似乎需要很多开销。

我发现的所有答案都会通过监督工作池而导致大量开销,我正在寻找一种解决方案,可以在工作人员级别自动获取终止信号。

【问题讨论】:

    标签: python multithreading exit flags pool


    【解决方案1】:

    我正在寻找一种能够在工作人员级别自动获取终止信号的解决方案。

    AFAIK,那不存在。 Pool 对象的方法(如Pool.terminate)应该在创建池的进程中使用。

    您可以使用Pool.imap_unordered。这会在结果上返回一个迭代器在父进程中,一旦结果可用,就会产生结果。只要弹出想要的结果,您就可以使用Pool.terminate()

    编辑

    • 通过查看 3.5 实现,starmap_async 返回一个 MapResult 实例,它不是一个迭代器。
    • 您可以将多个输入包装在一个元组中,并在其中的列表上使用 imap_unordered

    【讨论】:

    • 我目前正在使用“starmap_async”。它产生一个迭代器,但似乎在所有结果返回之前一直处于阻塞状态。 (它不应该)。这是一个已知的错误?我搜索并没有发现任何记录的问题。此外,“imap_unordered”是否支持多个函数输入?最后一点。在上面的代码中,'killpool()' 在工作池完成后立即运行,我只希望它在第一次返回时终止。
    • @DerrickCheek 我已经编辑了我的答案以解决您的评论。
    • 感谢您的更新!我回家后会试试这个。我一定混淆了迭代器的返回。我在尝试映射多个输入时遇到了几个问题,有些结果给出了迭代器,而有些则没有,哈哈。希望这有效。我会及时通知你。
    • @DerrickCheek Pool 方法的文档当然可以在这方面使用一些改进。因此,当有疑问时,请使用来源。 :-)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-07-09
    • 1970-01-01
    • 1970-01-01
    • 2016-05-31
    • 1970-01-01
    • 2019-08-01
    • 1970-01-01
    相关资源
    最近更新 更多