Spark如何读取目录名无分区名的分区Parquet数据集?
问题原因
你的判断完全正确,报错核心是Spark默认读取目录时会自动识别键=值格式的命名分区目录,你当前的路径层级只有纯数字分区值、没有分区名,Spark解析目录结构失败无法正常推断Parquet的统一Schema,所以抛出该错误。
更优解决方案
以下三个方案都不需要修改源存储结构,也不需要自行枚举所有底层文件路径:
方案1:开启递归文件查找(最便捷,Spark 3.0+支持)
开启recursiveFileLookup参数,让Spark忽略目录分区结构,直接递归扫描所有子目录下的Parquet文件,自动合并推断Schema:
df = spark.read \ .option("recursiveFileLookup", "true") \ .parquet('s3:\\my-bucket\files\14')
该方案无需提前定义Schema,代码量最少,适合快速开发场景。
方案2:手动指定Schema(最稳定,性能最优)
先读取单个文件拿到Schema,后续读取整个目录时直接传入指定的Schema,跳过Spark自动推断步骤:
# 先读取单个文件获取schema,也可以根据业务逻辑自行定义StructType single_file_df = spark.read.parquet('s3:\\my-bucket\files\14\09\12\file.pq') parquet_schema = single_file_df.schema # 读取上层目录时传入schema df = spark.read \ .schema(parquet_schema) \ .parquet('s3:\\my-bucket\files\14')
该方案避免了Spark扫描所有文件推断Schema的开销,生产环境优先推荐使用。
方案3:额外提取路径分区值
如果需要把路径里的日、月、时作为可查询的分区列使用,可以结合input_file_name()函数自行解析路径提取:
from pyspark.sql.functions import input_file_name, split, element_at df = spark.read \ .option("recursiveFileLookup", "true") \ .parquet('s3:\\my-bucket\files\14') \ .withColumn("file_path", input_file_name()) \ .withColumn("day", element_at(split("file_path", "/"), -3).cast("int")) \ .withColumn("month", element_at(split("file_path", "/"), -2).cast("int")) \ .withColumn("hour", element_at(split("file_path", "/"), -1).cast("int")) \ .drop("file_path")
注意根据你实际的S3路径层级调整索引值即可。
内容的提问来源于stack exchange,提问作者Umar.H
相关产品推荐
相关产品推荐

