如何按指定日期范围读取多文件夹Parquet数据到Spark DataFrame
问题场景
存储容器中有多个以日期为后缀的文件夹,示例路径:
dbfs:/mnt/input/raw/extract/pro_2023-01-01/parquet files here dbfs:/mnt/input/raw/extract/pro_2023-01-02/parquet files here dbfs:/mnt/input/raw/extract/pro_2023-01-03/parquet files here dbfs:/mnt/input/raw/extract/pro_2023-01-04/parquet files here dbfs:/mnt/input/raw/extract/pro_2023-01-05/parquet files here
单独读取单个文件夹Parquet数据到Spark DataFrame性能正常:
df = spark.read.parquet("dbfs:/mnt/input/raw/extract/pro_2023-01-05/")
但按周加载多日数据时,当前做法是先读取所有2023年前缀的文件夹,转换后用SQL过滤日期范围:
df = spark.read.parquet("dbfs:/mnt/input/raw/extract/pro_2023-*/") # 执行一些转换并添加新列 df.createOrReplaceTempView("Alldata")
select * from Alldata where cast(FolderDate as date) BETWEEN '2023-01-01' AND '2023-01-07'
这种方式会扫描所有文件夹,耗时极长,希望能在读取阶段直接拉取指定日期范围的数据。
解决方案
1. 生成指定日期范围的路径列表读取
先通过代码生成目标日期对应的文件夹路径,再批量读取,Spark只会扫描这些指定路径:
from datetime import datetime, timedelta start_date = datetime(2023, 1, 1) end_date = datetime(2023, 1, 7) date_list = [] current_date = start_date while current_date <= end_date: date_str = current_date.strftime("%Y-%m-%d") date_list.append(f"dbfs:/mnt/input/raw/extract/pro_{date_str}/") current_date += timedelta(days=1) # 读取指定路径的Parquet数据 df = spark.read.parquet(*date_list) # 后续执行转换操作
2. 调整文件夹结构为分区格式(推荐)
把文件夹命名改成Spark支持的分区格式,比如pro=2023-01-01,结构变为:
dbfs:/mnt/input/raw/extract/pro=2023-01-01/parquet files here dbfs:/mnt/input/raw/extract/pro=2023-01-02/parquet files here ...
读取时Spark会自动识别pro作为分区列,此时用where过滤日期范围,Spark会将过滤条件下推到读取阶段,只扫描符合条件的分区:
df = spark.read.parquet("dbfs:/mnt/input/raw/extract/") # 直接过滤分区列,无需全量扫描 df_filtered = df.where("pro between '2023-01-01' and '2023-01-07'") # 后续执行转换操作
这种方式不仅提升性能,还能让分区管理更规范,后续维护更方便。
3. 使用文件系统支持的范围通配符
如果你的存储系统(比如DBFS)支持bash风格的范围通配符,可以直接用通配符指定日期范围:
# 读取1月1日到7日的文件夹 df = spark.read.parquet("dbfs:/mnt/input/raw/extract/pro_2023-01-{01..07}/") # 后续执行转换操作
注意这种通配符的支持取决于底层存储系统,需要确认你的环境是否兼容。
内容的提问来源于stack exchange,提问作者user1403789
相关产品推荐
相关产品推荐

