为何Dask按索引字段筛选Parquet分区数据的速度远慢于直接读取?
优化方案
问题根因
你当前的操作流程会先加载全量9.5G的所有分区数据到Dask中,再进行索引设置和筛选,全量数据的加载、处理会产生大量不必要的IO和计算开销。
即使完成set_index后,你使用的df.loc[df.category == 'category-1']写法会触发全表逐行判断category字段的值,没有利用到索引的分区裁剪能力,依然会扫描所有分区的数据,所以速度远低于直接读取单个分区。
最优解决方法:读时分区下推
直接在read_parquet阶段传入过滤条件,Dask会自动识别过滤字段是分区列,仅加载符合条件的分区数据,完全不触碰其他分区的文件,性能和pandas直接读取对应分区文件夹基本一致。
示例代码:
# 直接下推过滤条件,无需提前加载全量数据 df = ddf.read_parquet( '/tmp/dasktest/', filters=[('category', '=', 'category-1')], npartitions=16 ) data_df = df.compute()
全量加载场景的优化
如果业务需要先加载全量数据做多次不同分类的筛选,可以调整筛选写法,利用已设置的索引特性加速:
你已经将category设置为Dask DataFrame的索引,无需再通过列比较的方式筛选,直接用索引取值即可,Dask会根据索引的分区元数据直接定位到对应数据块,避免全表扫描:
# 优化后写法,仅扫描对应分区的数据 data_df = df.loc['category-1'].compute()
额外优化建议
如果筛选后不需要用到所有字段,可在read_parquet时传入columns参数指定需要读取的列,进一步减少IO开销。
内容的提问来源于stack exchange,提问作者Jeff
相关产品推荐
相关产品推荐

