Dask Distributed如何通过client.submit选择性调度CPU/GPU worker
可运行实现示例
前置说明
调度的核心前提是worker启动时必须预先声明对应资源标签,资源名大小写敏感,提交任务时的资源名必须和worker声明的完全一致。
步骤1:启动带资源标签的Dask集群
你可以通过两种方式完成:
方式1:命令行手动启动
- 启动调度器:
dask scheduler - 启动CPU worker(示例配置4核CPU资源):
dask worker tcp://127.0.0.1:8786 --resources "CPU=4" - 启动GPU worker(示例配置1块GPU资源,多卡可设为对应卡数):
dask worker tcp://127.0.0.1:8786 --resources "GPU=1"
方式2:代码内直接创建集群(适合本地测试)
from dask.distributed import LocalCluster, Client # 集群配置:1个CPU worker(声明4个CPU资源)、1个GPU worker(声明1个GPU资源) cluster = LocalCluster( n_workers=2, resources=[{"CPU": 4}, {"GPU": 1}], threads_per_worker=2 ) client = Client(cluster) print(client.scheduler_info) # 可查看各worker的资源配置,确认标签生效
步骤2:使用client.submit按需调度任务
基础调度示例
# 调度到CPU worker执行的任务 def cpu_task(): import os return f"CPU task executed on worker: {os.uname()[1]}" # 调度到GPU worker执行的任务 def gpu_task(): import cudf import os df = cudf.DataFrame({"a": [1,2,3], "b": [4,5,6]}) return f"GPU task executed on worker: {os.uname()[1]}, cudf test passed: {len(df)=}" # 提交任务时指定resources参数,键为资源名,值为需要占用的资源数量 cpu_future = client.submit(cpu_task, resources={"CPU": 1}) gpu_future = client.submit(gpu_task, resources={"GPU": 1}) # 获取结果 print(client.gather(cpu_future)) print(client.gather(gpu_future))
dask-cudf + xgboost训练适配示例
如果是训练任务,需要将训练任务绑定到GPU worker,注意dask-xgboost的资源指定方式和client.submit一致:
import xgboost as xgb from dask_cudf import read_csv # 读取分布式数据 df = read_csv("s3://your-dataset-path/*.csv") X = df.drop("label", axis=1) y = df["label"] dtrain = xgb.dask.DaskDMatrix(client, X, y) # 指定训练任务使用GPU资源 params = { "tree_method": "gpu_hist", "objective": "binary:logistic", "max_depth": 8 } # 提交训练任务时指定resources参数 train_future = client.submit( xgb.dask.train, client, params, dtrain, num_boost_round=100, resources={"GPU": 1} ) booster = client.gather(train_future)["booster"]
常见问题排查
- 检查资源名大小写是否完全匹配,
GPU和gpu会被识别为两种不同的资源 - 确认worker启动时确实配置了对应资源,可通过Dask仪表盘的
Workers标签页查看各worker的资源配置 - 如果任务一直处于pending状态,说明当前没有匹配对应资源要求的空闲worker,可扩容对应资源的worker实例
内容的提问来源于stack exchange,提问作者Ehsan Fathi
相关产品推荐
相关产品推荐

