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

Dask分布式:Worker端初始化全局有状态参数及复制异常问题

本地Dask集群Worker端预加载全局有状态参数的问题

问题背景

我正在本地搭建Dask集群(调度器与Worker均部署在localhost),集群创建代码如下:

cluster = SSHCluster(["localhost", "localhost"],
                    connect_options={"known_hosts": None},
                    worker_options={"n_workers": params["n_workers"], },
                    scheduler_options={"port": 0, "dashboard_address": ":8797"},)

client = Client(cluster)

我希望在Worker端初始化有状态全局参数,供后续分配给Worker的任意Worker方法调用。

尝试方案与现存问题

我尝试使用client.register_worker_plugin方法实现,代码如下:

def read_only_data(jsonfilepath):
    with open(jsonfilepath, "r") as readfile:
        return json.load(readfile)

def main():
    cluster = SSHCluster(params) # simplified
    client = Client(cluster)
    plugin = read_only_data(jsonfilepath)
    client.register_worker_plugin(plugin, name="read-only-data")

但发现只读数据是在客户端初始化后复制到Worker端的,大数据场景下会产生额外通信开销。既然Worker端本身就有目标JSON文件,每次调用Worker方法都重复加载又会造成冗余执行,我想了解是否存在Worker创建完成、任务分配前的初始化事件,可通过该事件预加载数据结构作为全局有状态参数,供后续Worker方法共享使用。

更新:尝试replicate方法后的异常情况

按照建议使用client.replicate()后,处理280MB的JSON文件时,该操作卡住超过1小时(单进程加载该文件仅需20秒)。Dask Dashboard显示每个Worker内存占用约2GB且持续增长,同时存在网络活动,最终脚本抛出错误:

distributed.comm.core.CommClosedError: in <TCP (closed) ConnectionPool.replicate local=tcp://192.168.0.6:38498 remote=tcp://192.168.0.6:46171>: TimeoutError: [Errno 10] Connection timed out

内存占用远超预期(280MB×6≈1.7GB),虽然Dask文档说明replicate会将数据复制到Worker,但无法解释此次内存异常与复制耗时过长的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 01:42:39