如何让Dask worker执行完任务后自动终止,无需关停整个集群
解决方案
以下两种方式都可以满足你的需求,无需关停整个集群,也不会影响任务返回值回传:
方案1:配置worker单任务后自动退役(适配你的场景,无需改任务代码)
你的worker本身是单线程配置,直接设置worker执行1个任务后自动优雅退出即可,Dask调度器会自动拉起新的worker进程替换,全程不会影响其他运行中的任务:
- 命令行启动worker示例:
dask worker <调度器地址> --nworkers <进程数> --nthreads 1 --maximum-tasks 1 - Python代码初始化集群示例:
from dask.distributed import LocalCluster cluster = LocalCluster( n_workers=4, threads_per_worker=1, worker_kwargs={"maximum_tasks": 1} )
该方案下Dask会自动等待任务返回值完全回传到调度器后,再终止旧worker进程,完全规避内存泄漏问题。
方案2:特定任务完成后主动触发worker退役
如果你只需要让调用了泄漏C库的任务触发worker重启,其他普通任务可以复用worker,可以在对应任务的末尾加入主动退役逻辑:
from distributed import get_worker def your_leaky_task(*args): # 执行业务逻辑,调用存在内存泄漏的C库 res = your_c_lib_call() # 任务执行完成后触发当前worker优雅退役 worker = get_worker() worker.retire() # 正常返回结果即可,Dask会先完成结果回传再终止worker return res
注意:两种方案都默认依赖Dask的worker自动重启机制,如果你通过官方默认的
dask worker命令或者内置集群管理器启动集群,该机制默认生效,无需额外配置。
内容的提问来源于stack exchange,提问作者Throw Away Account
相关产品推荐
相关产品推荐

