如何并行计算字典中存储的多个Dask DataFrame实例?
如何并行计算字典中的Dask DataFrame?
你遇到的问题是因为dask.compute(*container_dict)实际解包的是字典的键而非值,所以无法触发对应Dask DataFrame的并行计算。只需调整为解包字典的值,再将计算结果映射回原字典即可实现一次并行计算。
解决方案代码
import dask.dataframe as dd container_dict = {} for index, value in enumerate(comb_dict_stock): container_dict[index] = ddf.loc[index] # 并行计算所有字典中的Dask对象 computed_results = dask.compute(*container_dict.values()) # 将计算结果重新关联到原字典的键 container_dict = dict(zip(container_dict.keys(), computed_results))
更简洁的写法
可以将结果映射步骤合并为一行:
container_dict = dict(zip(container_dict.keys(), dask.compute(*container_dict.values())))
原理说明
dask.compute(*container_dict.values())会一次性提交所有字典内的Dask任务,利用Dask的并行调度器同时执行,避免了循环调用compute()带来的串行执行开销。- 用
zip将原字典的键和计算后的结果一一对应,重建完成计算的字典。
内容的提问来源于stack exchange,提问作者user4933
相关产品推荐
相关产品推荐

