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

