【发布时间】: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