如何确定dask.distributed.cluster可支持的最大worker数量
问题
如何获取dask.distributed.client可访问的最大worker数量,要求包含尚未启动的worker,同时兼容自适应和非自适应两种扩缩容策略?
我开发了一款用于函数自适应并行执行的库,可提前规划待执行的点位。为此我们需要获知可并行访问的worker总数,目前我们使用client.ncores()函数来获取该数值。
但该方案存在多个问题:
- 首先需要worker处于运行状态,因为
ncores仅返回当前实时的运行数据。 - 若集群启用了自适应扩缩容,我们需要获知可获取的最大worker数量而非当前运行数量。
因此我想了解是否有程序化的方法可以通过client查询,确定dask集群可扩容到的最大worker数量。
解决方案
你可以通过 Dask 客户端的集群属性与调度器元数据,统一查询兼容自适应、非自适应模式的最大可扩容 worker 数量,无需依赖当前已运行的 worker 状态。
通用实现代码(适配绝大多数本地/托管集群场景)
from distributed import Client def get_cluster_max_workers(client: Client) -> int: # 优先读取自适应扩缩容的最大配置 if hasattr(client.cluster, "adapt") and client.cluster.adapt is not None: max_workers = client.cluster.adapt.maximum return max_workers if max_workers is not None else 0 # 非自适应固定规模集群读取预配置的总worker数 else: return len(getattr(client.cluster, "worker_spec", []))
场景说明
- 自适应扩缩容集群:
client.cluster.adapt.maximum直接返回集群配置的 worker 数量上限,不受当前已启动 worker 数量影响 - 非自适应固定规模集群:
client.cluster.worker_spec存储了集群预配置的所有 worker 规格,取长度即可得到最大可启动 worker 总数 - 若你需要获取最大核心总数而非 worker 数量,可以遍历
client.cluster.worker_spec中每个 worker 配置的核心数求和即可。
远程独立集群兼容方案
如果连接的是独立部署的远程 Dask 集群,客户端没有挂载 cluster 属性,可以通过调度器元数据接口查询:
def get_remote_cluster_max_workers(client: Client) -> int: scheduler_info = client.scheduler_info() adapt_config = scheduler_info.get("extensions", {}).get("adaptive", {}) # 自适应场景取配置上限,非自适应场景取预配置worker总数 if "maximum" in adapt_config: return adapt_config["maximum"] return len(scheduler_info.get("workers", {})) if not adapt_config else 0
内容的提问来源于stack exchange,提问作者Anton Akhmerov
相关产品推荐
相关产品推荐

