Dask Gateway环境下配置worker自定义资源的实现方法咨询
方案1:通过Dask Gateway Helm Chart配置(支持用户自定义选项)
- 管理员全局默认配置:在Helm的values.yaml中添加如下配置,可给所有集群的Worker默认设置固定资源
gateway: clusterConfig: worker: extraArgs: - --resources - task_limit=1
- 支持用户自定义配置:如果需要让用户创建集群时可自行调整资源参数,需要先在Helm配置中新增集群选项字段,再将参数映射到Worker启动参数:
gateway: clusterOptions: - name: worker_task_limit type: integer default: 1 min: 1 label: 单Worker最大并行任务数 clusterConfig: worker: extraArgs: - --resources - task_limit={{ cluster_options.worker_task_limit }}
用户侧创建集群时即可传入自定义参数:
from dask_gateway import Gateway gateway = Gateway() # 自定义设置单Worker并行任务上限为2 cluster = gateway.new_cluster(worker_task_limit=2)
方案2:Worker启动后动态配置
你之前的代码不生效是因为Worker本地修改资源后不会主动同步到Scheduler,因此Scheduler侧存储的Worker资源信息不会更新,可使用如下方法完成同步更新:
def update_worker_resources(dask_worker): # 更新Worker本地资源记录 dask_worker.set_resources(task_limit=1) # 主动将资源配置同步到Scheduler dask_worker.scheduler.update_resources( worker=dask_worker.address, resources=dask_worker.total_resources ) return dask_worker.total_resources # 所有Worker执行资源更新操作 client.run(update_worker_resources)
如果仅需要给指定Worker设置资源,也可以直接调用Scheduler接口完成更新:
# 替换为目标Worker的实际地址 worker_addr = "tls://255.0.92.232:39305" client.scheduler.update_resources(worker=worker_addr, resources={"task_limit": 1})
注意:配置完Worker资源后,提交任务时需要指定对应的资源参数,调度器才会按照资源限制分配任务:
# 示例:提交任务时声明需要占用1个task_limit资源 future = client.submit(your_function, *args, resources={"task_limit": 1})
内容的提问来源于stack exchange,提问作者bill_e
相关产品推荐
相关产品推荐

