如何在分布式worker中执行Dask Graph中的每个函数?
要让你写的Dask Graph在分布式worker上执行,只需要调整执行入口即可,原有Graph结构不需要改动,确实需要用分布式场景下的get方法替代本地执行的get,具体操作如下:
- 首先确保你已经安装了dask分布式依赖:
pip install dask distributed - 初始化Dask客户端,连接到你的分布式集群。如果只是本地测试分布式调度逻辑,也可以直接启动本地临时集群。
修改后的完整可运行代码如下:
from dask.distributed import Client # 初始化客户端,连接集群时填入你实际的调度器地址,不传参数会默认启动本地集群做测试 client = Client("tcp://your-scheduler-ip:8786") data1 = [1000, 2000, 3000, 4000, 5000 ] data2 = [ 100, 200, 300, 400, 500 ] def sum_all(x): return sum(x[0]) + sum(x[1]) def sum_10(x): return x + 10 graph2 = { 'data' : (data1, data2) } graph2['out2'] = (sum_10, 'out1') graph2['out1'] = (sum_all, 'data') # 替换原来的本地get为client.get即可,任务会自动分发到worker执行 result = client.get(graph2, ['out1', 'out2']) print(str(result)) # 输出结果和本地运行一致:(16500, 16510)
调用client.get的过程中,客户端会自动将你定义的函数、任务图序列化后提交到调度器,调度器会拆分任务分配到各个worker执行,最后把结果汇总返回给客户端,你不需要额外做其他任务提交操作。如果后续要提交更大规模的任务,只需保证自定义函数可被序列化、worker节点的运行环境和客户端环境一致即可。
内容的提问来源于stack exchange,提问作者ps0604
相关产品推荐
相关产品推荐

