如何检查运行中的Dask调度器?启动指定Worker数本地集群前的校验
检查现有Dask调度器并启动指定Worker的本地集群
要检查是否存在通过dask-scheduler命令启动的本地集群,最简单的方式就是尝试连接默认的调度器地址(tcp://127.0.0.1:8786,这是dask-scheduler的默认监听地址)。如果连接成功,说明已有集群在运行;如果连接被拒绝,再启动你自定义的本地集群即可。
我通常会用try-except块来优雅处理这两种情况,给你一个现成的代码实现:
from dask.distributed import Client, LocalCluster # 默认的dask-scheduler监听地址 default_scheduler_addr = "tcp://127.0.0.1:8786" try: # 先尝试连接现有调度器 client = Client(default_scheduler_addr) print(f"成功连接到已运行的Dask调度器:{default_scheduler_addr}") print(f"当前集群的Worker数量:{len(client.scheduler_info()['workers'])}") except ConnectionRefusedError: # 没有检测到现有调度器,启动自定义本地集群 print("未发现运行中的Dask调度器,启动新的本地集群...") cluster = LocalCluster(n_workers=8, ip='127.0.0.1') client = Client(cluster) print(f"新集群已启动,调度器地址:{client.scheduler_info()['address']}") # 后续就可以用client提交任务了 # 示例:result = client.submit(lambda x: x**2, 10)
额外说明:
- 如果你启动
dask-scheduler时修改了默认端口(比如用dask-scheduler --port 9999),记得把default_scheduler_addr里的端口改成对应数值。 - 要是想更严谨,还可以通过
client.scheduler_info()返回的信息,进一步确认调度器的运行状态,避免连接到僵死进程。 - 这种方式既能复用现有集群资源,也能在无集群时自动启动新集群,适配你的需求。
内容的提问来源于stack exchange,提问作者medRa
相关产品推荐
相关产品推荐

