如何实现Dask两个独立集群间的跨集群数据拉取
跨Dask集群数据传输方案
两个独立Dask LocalCluster之间传输数据完全可行,根据数据量大小可以选以下两种落地性强的方案:
方案1:内存直连传输(适合小批量、GB级以下数据)
- 第一步:获取第一个集群(clus1)的调度器地址
在第一个Jupyter Notebook中执行print(cluster.scheduler_address),会输出类似tcp://127.0.0.1:37245格式的地址,记录这个地址备用。 - 第二步:在第二个Notebook中同时初始化两个集群的客户端连接
你之前在第二个Notebook里只连接了clus2,现在额外加一个指向clus1的客户端实例即可,示例代码:
from dask.distributed import Client # 连接当前Notebook对应的第二个集群,地址替换为clus2的scheduler_address client_clus2 = Client("tcp://127.0.0.1:替换为clus2的调度器端口") # 连接第一个集群,地址替换为第一步记录的clus1的scheduler地址 client_clus1 = Client("tcp://127.0.0.1:替换为clus1的调度器端口")
- 第三步:在第一个集群中发布需要传输的数据集
回到第一个Notebook,把你预处理完成的数据集发布到clus1的可查询目录中,让其他客户端可以拉取:
# df_processed替换为你实际完成预处理的pandas/dask对象 client_clus1.publish_dataset(processed_data=df_processed)
- 第四步:在第二个Notebook中拉取数据并提交到clus2计算
# 从clus1拉取已发布的数据集 df_trans = client_clus1.get_dataset("processed_data") # 将数据提交到clus2的计算链路,后续可以直接接你的处理逻辑 clus2_future = client_clus2.submit(lambda x: x, df_trans) # 拿到clus2上的可用数据对象 df_on_clus2 = clus2_future.result()
方案2:共享存储中转(适合GB级以上大数据量,稳定性最高)
如果数据量较大,内存直传容易触发OOM或者网络传输超时,优先选这种方案:
- 第一步:在第一个Notebook中把预处理完成的数据写入本地可共享访问的磁盘路径,推荐用列式压缩格式减少IO开销:
# pandas对象直接写parquet即可 df_processed.to_parquet("./dask_shared/processed_data.parquet")
- 第二步:在第二个Notebook中直接从上述共享路径读取数据,加载到clus2的计算流程中即可,不需要做额外的跨集群通信配置。
注意事项
- 同一台机器上启动多个LocalCluster时,必须保证scheduler、dashboard、worker端口互不冲突,你当前配置的dashboard端口8789、8790不存在冲突,只要启动时没有端口占用报错即可正常使用。
- 内存直传模式下,待传输的数据必须提前通过
publish_dataset注册到集群调度器,仅存在于Notebook内核内存中的对象无法被其他集群直接访问。 - 如果两个集群部署在不同物理机上,只要网络互通,上述内存直连方案同样生效,只需要把scheduler地址替换为对应机器的IP+端口即可。
内容的提问来源于stack exchange,提问作者XGB
相关产品推荐
相关产品推荐

