You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何编程创建多机器Dask远程Worker并复用Client

解决方案:远程Dask Worker部署与Client复用

问题1:在指定远程机器创建Worker

你当前用distributed.Worker类启动Worker的方式,本质是在本地进程中创建Worker——哪怕指定了host参数也只是设置Worker对外暴露的地址,不会帮你把Worker部署到远程机器上。要实现多机器远程部署,得用Dask专门的远程集群工具:dask_ssh(基于SSH连接远程机器启动Worker)。

步骤与代码示例

  1. 先安装依赖:
pip install dask-ssh
  1. 用SSHCluster连接远程Scheduler,并在目标机器启动Worker:
from dask_ssh import SSHCluster
from distributed import Client

# 远程Scheduler地址
scheduler_address = "x.x.x.x:8786"
# 要部署Worker的目标机器列表(确保当前机器能无密码SSH登录这些机器)
worker_hosts = ["ip1", "ip2", "ip3"]

# 创建SSHCluster,关联到已有的远程Scheduler
cluster = SSHCluster(
    worker_hosts,
    scheduler_address=scheduler_address,
    # 可选:自定义Worker配置,比如进程数、线程数
    worker_options={"nprocs": 2, "nthreads": 4}
)

# 确认Worker已成功连接到Scheduler
print("Worker部署完成,当前集群状态:", cluster.scheduler_info())

这样就能在ip1、ip2、ip3这几台机器上启动Worker,并且自动关联到你指定的远程Scheduler了。


问题2:跨async函数复用Client并管理Worker生命周期

你之前的代码问题在于:async with上下文管理器会在代码块退出时自动关闭Client和Worker,所以你return的Client在上下文结束后已经是关闭状态了。要复用Client,我们需要把Worker的生命周期管理和Client的使用分开,同时确保最后能正确清理资源。

方案1:同步Client复用(适合大多数场景)

如果你的业务代码是同步的,直接用同步Client即可,不需要异步上下文:

# 承接上面的cluster实例
client = Client(cluster)

# 在任意地方复用这个client
def custom_function(client):
    # 执行Dask任务
    future = client.submit(lambda x: x + 1, 10)
    result = future.result()
    print("任务结果:", result)
    # 也可以调用自定义模块的函数
    # custom_module.call_some_fn_to_use_dask_client(client)

# 调用自定义函数
custom_function(client)

# 所有任务完成后,关闭Client和Worker集群
client.close()
cluster.close()

方案2:异步Client复用(如果必须用异步代码)

如果你的业务逻辑是异步的,可以手动管理Client和Worker的生命周期,避免上下文提前关闭:

import asyncio
from dask_ssh import SSHCluster
from distributed import Client

scheduler_address = "x.x.x.x:8786"
worker_hosts = ["ip1", "ip2", "ip3"]

async def init_cluster_and_client():
    # 创建异步模式的SSHCluster
    cluster = await SSHCluster(
        worker_hosts,
        scheduler_address=scheduler_address,
        asynchronous=True
    )
    # 创建异步Client
    client = await Client(cluster, asynchronous=True)
    return cluster, client

async def use_client(client):
    # 异步执行任务
    future = client.submit(lambda x: x + 1, 10)
    result = await future
    print("异步任务结果:", result)
    # 调用自定义异步函数
    # await custom_module.async_call_some_fn(client)

async def main():
    cluster, client = await init_cluster_and_client()
    try:
        await use_client(client)
        # 在这里可以多次复用client执行其他任务
    finally:
        # 确保最后清理资源,避免远程Worker残留
        await client.close()
        await cluster.close()

# 运行主异步函数
asyncio.run(main())

关键注意点

  • 如果你需要在同步代码中调用异步Client,可以用client.sync()方法包装异步操作,比如:client.sync(use_client, client)
  • 一定要确保在所有任务完成后关闭cluster和client,否则远程Worker会一直运行占用资源。

内容的提问来源于stack exchange,提问作者spiralarchitect

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.06 10:07:33