Spark导出分区Parquet后,Dask(pyarrow引擎)无法读取分区列求助
问题分析与解决方案
这个问题的核心在于Spark写入分区Parquet的方式和Dask+PyArrow读取时的分区解析逻辑不匹配:
Spark按bhello分区写入Parquet时,会把分区列作为Hive风格的目录结构(比如bhello=hello/、bhello=Yo/这样的文件夹)存储,而不是将分区列数据嵌入到Parquet数据文件中。但Dask使用PyArrow引擎读取时,默认不会自动解析这种目录中的分区信息,所以bhello列不会出现在读取后的DataFrame里。
修复方法:指定Hive风格分区解析
只需要在dd.read_parquet中添加partitioning='hive'参数,告诉PyArrow去识别目录中的Hive格式分区键:
df2 = dd.read_parquet( 'hdfs://127.0.0.1:8020/tmp/test/outputParquet10', engine='pyarrow', partitioning='hive' # 关键参数,开启Hive分区解析 )
执行这段代码后,bhello列就会被正确加载到df2中,你的断言也会通过。
额外验证点
- 检查目录结构:确认HDFS上的输出目录下确实存在
bhello=xxx格式的子文件夹,确保Spark正确完成了分区写入。 - 版本兼容性:确保你的Dask(>=2021.06.0)和PyArrow(>=3.0.0)版本足够新,旧版本对Hive分区的支持可能存在缺陷。
替代方案:用PyArrow Dataset读取
如果上述方法不生效,也可以直接使用PyArrow Dataset API构建数据集,再转换为Dask DataFrame:
import pyarrow.dataset as ds # 构建PyArrow数据集,指定Hive分区解析 dataset = ds.dataset( 'hdfs://127.0.0.1:8020/tmp/test/outputParquet10', format='parquet', partitioning='hive' ) # 转换为Dask DataFrame df2 = dd.from_pyarrow_dataset(dataset)
内容的提问来源于stack exchange,提问作者pranav kohli
相关产品推荐
相关产品推荐

