如何修改Dask本地集群Worker名称并精准识别新增Worker
我所咨询的是本地集群场景下的相关问题。
以下是用于获取Worker地址与名称的实现代码:
def f(): worker = get_worker().name return worker client.run(f)
运行输出结果:
{'tcp://127.0.0.1:58709': 0, 'tcp://127.0.0.1:58710': 2, 'tcp://127.0.0.1:58711': 1}
返回结果为字典结构,键为Worker的访问地址,值为对应Worker的名称。我想了解是否可以在创建Client时直接指定Worker名称,或是存在其他可行的替代方案?
需要说明的是,仅修改上述返回字典的取值无法更改Worker的实际名称。
需求背景
我的目标是首先创建1个Worker,将数据预处理任务分配至该节点运行;预处理完成后构建逻辑回归、随机森林分类器两个模型,随后通过scale方法扩容Worker数量,将两个模型的训练任务分别分配至两个新增Worker上运行。
目前遇到的问题是无法识别哪些是扩容产生的新增Worker:默认Worker名称为worker_0、worker_1、worker_2这类格式,我无法判断扩容后新增的Worker对应哪个标识。我推测新增Worker的名称是按增量规则生成的,但缺乏有效依据验证该猜想,因此考虑通过自定义Worker名称的方式更便捷地追踪Worker状态,实现任务的定向调度。
解决方案
自定义Worker名称实现方式
创建Client时无法直接指定Worker名称,需要在初始化集群实例时传入Worker配置,扩容阶段也可以手动指定新增Worker的名称,从根源上避免序号识别问题:
from dask.distributed import Client, LocalCluster # 初始化集群,启动1个预处理专用节点 cluster = LocalCluster(n_workers=1, worker_kwargs={"name": "preprocess_node"}) client = Client(cluster) # 预处理任务执行完成后,扩容时指定两个训练节点的自定义名称 cluster.scale(3, workers=["preprocess_node", "lr_train_node", "rf_train_node"])
执行client.run(lambda: get_worker().name)即可看到每个Worker地址对应绑定的自定义名称,后续提交任务时直接用名称指定运行节点即可。
无自定义名称的替代识别方案
如果不需要固定Worker名称,可以通过扩容前后的Worker列表差集识别新增节点,完全不依赖名称生成规则:
# 预处理阶段记录初始Worker地址集合 init_worker_set = set(client.scheduler_info()["workers"].keys()) # 触发扩容并等待所有节点就绪 cluster.scale(3) client.wait_for_workers(3) # 计算差集得到新增的两个Worker地址 current_worker_set = set(client.scheduler_info()["workers"].keys()) new_workers = list(current_worker_set - init_worker_set) lr_worker_addr, rf_worker_addr = new_workers
后续提交任务时通过workers参数传入对应地址即可定向调度,例如client.submit(train_lr_model, train_data, workers=lr_worker_addr)。
默认名称规则说明
本地集群默认Worker名称确实从0开始按启动顺序递增生成,但如果运行过程中有Worker异常退出、释放占用的序号,后续新增Worker会优先复用空缺序号,直接通过序号判断节点新旧存在逻辑风险,不推荐在生产流程中使用该逻辑。
内容的提问来源于stack exchange,提问作者XGB

