如何确保Dask查询分区DataFrame时仅读取必要磁盘文件
如何确保Dask仅读取匹配过滤条件的目标Parquet文件
你当前使用的Hive风格分区存储格式,本身就是Dask和PyArrow原生支持的分区裁剪场景,只要注意几个要点,就能保证只读取目标文件,不会加载全量数据:
- 最稳妥的方式是在读取阶段就传入过滤条件,直接在构建任务图的时候把不符合条件的分区全部排除,从根源上避免无关文件被读取,示例代码如下:
import dask.dataframe as dd # 读取时直接指定过滤规则,仅加载匹配分区的文件 df = dd.read_parquet( "dataset/", engine="pyarrow", filters=[[("field1", "==", "A"), ("field2", "==", "X")]] )
这种写法下,后续所有计算都只会访问dataset/field1=A/field2=X/data.parquet这一个文件,不会有额外IO。
- 如果是像你写的那样先读全量数据集、再用
query()做过滤,只要满足两个条件,Dask也会自动把过滤条件下推到读取层,自动完成分区裁剪:- 过滤操作前没有做过会打乱分区结构的操作,比如重分区、shuffle类操作(join、全量groupby聚合等)
- 过滤条件直接作用在分区列上,没有对分区列做函数转换(比如
field1.str.lower() == 'a'这类写法就无法被下推,会导致全量文件被读取)
你写的df.query("field1 == 'A' and field2 == 'X'")是符合下推要求的,默认就只会读取目标文件。
- 快速验证裁剪是否生效:可以在过滤后打印
filtered_df.npartitions,如果返回值是1,就说明已经裁剪到只剩目标分区对应的文件了,没有加载其他3个分区。
如何查看Dask实际读取了哪些文件
有几个可直接落地的方法:
- 方法1:开启PyArrow数据集模块的DEBUG日志,执行计算时会直接在控制台打印所有被访问的文件路径,示例代码:
import logging import sys # 配置日志输出到控制台 logging.basicConfig(stream=sys.stdout, level=logging.WARNING) # 单独开启pyarrow数据集模块的debug日志,会输出所有文件访问记录 logging.getLogger("pyarrow.dataset").setLevel(logging.DEBUG) # 执行计算,控制台会打印所有被打开读取的Parquet文件路径 res = df.query("field1 == 'A' and field2 == 'X'").compute()
- 方法2:如果使用Dask Distributed集群,可以打开Dask Dashboard的Worker日志页,每个Worker读取文件的记录都会实时打印在日志里;如果是本地调试不想开集群,直接用上面的日志方法即可。
- 额外技巧:不需要执行计算也能判断哪些文件会被读取,过滤后打印DataFrame的分区数,如果分区数和你预期匹配(比如这个场景下应该是1),说明裁剪已经生效,不会读取无关文件。
内容的提问来源于stack exchange,提问作者william_grisaitis
相关产品推荐
相关产品推荐

