【问题标题】:Terminating dask workers after jobs are done工作完成后终止工作人员
【发布时间】:2023-03-25 16:27:02
【问题描述】:

我正在尝试在集群上使用 dask,并且我有兴趣在所有工作完成后立即终止所有工作人员。 我试图用 retire_workers 方法做到这一点,但这似乎并没有杀死工人。 这是一个例子。

import time
import os
from dask.distributed import Client

def long_func(x):
    time.sleep(2)
    return 1

if __name__ == '__main__':
    C = Client(scheduler_file='sched.json')
    res = []
    for _ in range(10):
        res.append(C.submit(long_func, _))
    for r in res:
        r.result()
    workers = list(C.scheduler_info()['workers'])
    # C.run(lambda: os._exit(0), workers=workers)
    C.retire_workers(workers=workers, close_workers=True)

调度程序和工作人员使用以下命令启动:

dask-scheduler --scheduler-file sched.json 
dask-worker --scheduler-file sched.json --nthreads=1 --lifetime='5minutes'

希望在执行上面的 python 代码后,worker 会终止(20 秒后),但它不会,停留整整 5 分钟。有什么建议可以解决这个问题吗?

【问题讨论】:

    标签: python parallel-processing cluster-computing dask


    【解决方案1】:

    这将关闭连接的调度程序并退休工人:

    C.shutdown()
    

    【讨论】:

    • 您可能还想在此之后运行C.close() 以避免收到有关缺少调度程序的消息。
    【解决方案2】:

    我建议使用上下文管理器来管理集群 - 它既美观又干净。在本地工作时,当 RAM 内存被用尽并停止计算机时,我遇到了问题,但这里是我经常使用的示例:

    # start our Dask cluster
    from dask.distributed import Client,LocalCluster
    
    if __name__ == '__main__':
        cluster = LocalCluster()
        
        with Client(cluster) as client:
            print("scheduler host: ", client.scheduler.address)
            # do some stuff
    
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2020-01-04
      • 2015-06-14
      • 1970-01-01
      • 1970-01-01
      • 2021-01-11
      • 2014-03-12
      相关资源
      最近更新 更多