为何Dask Bags运行速度远低于串行计算?有哪些优化方案?
问题原因
- 调度与通信开销远大于计算收益:当前用到的单元素计算逻辑
g(x)运算量极低,而Dask执行并行计算时需要完成任务拆分、数据序列化跨进程传输、任务调度、结果汇总回收等一系列额外操作,这些操作的总耗时已经远远超过了并行计算节约的时间,反而比直接串行运行慢。 - 工具选型不匹配:
dask.bag是设计用来处理半结构化/非结构化数据(如文本、JSON、日志条目等)的模块,用来处理规整的Numpy数值数组会引入大量不必要的类型转换、遍历开销,本身运行效率就很低。 - 任务粒度过细:直接对单元素执行
map操作,相当于为每个元素都生成了一个微型任务,进一步放大了调度开销。
优化方案
方案1:仍使用dask.bag的优化方式
调整任务粒度,不要对单个元素做映射,改为对整个分区做向量化计算,大幅降低调度开销:
# 当前示例的g函数天然支持numpy数组输入,无需额外修改 def g(x): return np.sqrt(np.abs(x)) ** np.log(np.abs(x)) %%time # 并行计算 # 分区数建议和worker数保持1:1到4:1的比例即可,避免过多分区增加调度开销 b = db.from_sequence(test_array, npartitions=8) # 对每个分区整体做映射,直接向量化计算整个分区的数据 b = b.map_partitions(lambda part: g(np.array(part))) results_parallel = b.compute()
方案2:使用更适配数组场景的Dask Array(更推荐)
Dask Array是专门为大型多维数组设计的模块,完全兼容Numpy接口,序列化和调度开销远低于dask.bag,代码改动量极小:
import dask.array as da %%time # 把numpy数组转为Dask Array,自动分块 dask_arr = da.from_array(test_array, chunks="auto") # 直接调用原g函数即可,Dask会自动对每个分块做向量化计算 results_parallel = g(dask_arr).compute()
额外优化建议
- 如果实际业务中的函数逻辑无法向量化,属于纯单元素复杂计算,可以用
numba.jit先编译函数,再配合Dask运行,能同时提升单任务计算速度和降低整体开销。 - 避免设置远大于worker数量的分区数,过多的分区只会徒增调度成本。
- 数据规模更大的场景下,可以先把数组存为Zarr格式的本地文件,再用Dask加载计算,省去进程间传输原始数组的开销。
内容的提问来源于stack exchange,提问作者Whyjay Zheng
相关产品推荐
相关产品推荐

