如何将Parquet每个行组读取为Dask DataFrame的独立分区?
让Dask DataFrame按Parquet行组创建独立分区
你遇到的这个情况其实是Dask默认行为导致的——它并不会总是自动把每个Parquet行组拆成单独分区,主要是为了避免生成过多小分区带来的调度性能损耗。不过完全不用把数据拆成多个文件就能实现每个行组对应一个分区,下面给你两种可行的方法:
方法一:使用split_row_groups参数(推荐,适用于较新版本Dask)
从Dask 2021.06.0版本开始,read_parquet新增了split_row_groups参数,设置为True就能强制让每个行组成为一个独立分区:
import dask.dataframe as dd # 读取时开启行组分拆 df = dd.read_parquet("/tmp/test2.parquet", split_row_groups=True) print(df.npartitions) # 此时应该输出10,和Parquet文件的行组数量一致
这个参数会告诉Dask不要合并行组,直接将每个行组映射为一个分区,简单高效。
方法二:手动遍历行组读取(适用于旧版本Dask)
如果你使用的是更早的Dask版本,没有split_row_groups参数,可以借助PyArrow先获取行组数量,然后逐个读取每个行组再合并:
import dask.dataframe as dd import pyarrow.parquet as pq # 先获取Parquet文件的行组信息 pq_file = pq.ParquetFile("/tmp/test2.parquet") row_group_count = pq_file.num_row_groups # 逐个读取每个行组,存入列表 partition_dfs = [] for row_group_idx in range(row_group_count): partition_df = dd.read_parquet("/tmp/test2.parquet", row_groups=[row_group_idx]) partition_dfs.append(partition_df) # 合并所有分区 df = dd.concat(partition_dfs) print(df.npartitions) # 输出10
补充说明
Dask默认不拆分单个文件的行组,核心原因是分区大小的平衡:如果行组太小,创建大量小分区会增加任务调度的开销,反而降低整体性能。但如果你的行组大小足够合理(比如每个行组几十MB到几百MB),用上面的方法拆分就完全没问题,不需要拆分文件。
内容的提问来源于stack exchange,提问作者gerrit
相关产品推荐
相关产品推荐

