如何在PySpark中按文件名日期批量读取S3的Parquet文件
批量读取指定日期的Parquet文件到PySpark DataFrame
假设S3存储桶中有如下Parquet文件:
s3://data-raw-dev/GoogleAds/ad_group/cron_name=customer2/year=2024/month=02/googleads_ad_group_customer2_2024-02-20.parquet s3://data-raw-dev/GoogleAds/ad_group/cron_name=customer2/year=2024/month=02/googleads_ad_group_customer2_2024-02-19.parquet s3://data-raw-dev/GoogleAds/ad_group/cron_name=customer3/year=2024/month=02/googleads_ad_group_customer3_2024-02-20.parquet s3://data-raw-dev/GoogleAds/ad_group/cron_name=customer3/year=2024/month=02/googleads_ad_group_customer3_2024-02-19.parquet
读取单个Parquet文件到PySpark DataFrame的方法很简单:
df_staging = spark.read.parquet(s3_path) df_staging.show()
现需根据文件名中的日期(例如2024-02-19)批量读取多个文件到PySpark DataFrame,无需遍历客户名称逐个读取,有几种简洁的实现方式:
方法1:使用路径通配符匹配
利用*通配符匹配任意客户名称,同时指定目标日期的文件名后缀,直接构造匹配路径:
target_date = "2024-02-19" # 用*匹配所有客户,同时锁定目标日期的文件 s3_path_pattern = f"s3://data-raw-dev/GoogleAds/ad_group/cron_name=*/year=2024/month=02/*{target_date}.parquet" df_staging = spark.read.parquet(s3_path_pattern) df_staging.show()
方法2:使用pathGlobFilter参数过滤文件名
通过pathGlobFilter选项指定文件名的匹配规则,结合递归路径读取所有客户下的文件:
target_date = "2024-02-19" df_staging = spark.read.option("pathGlobFilter", f"*{target_date}.parquet") \ .parquet("s3://data-raw-dev/GoogleAds/ad_group/cron_name=*/year=2024/month=02/") df_staging.show()
方法3:读取多个指定日期的文件
如果需要一次性读取多个日期的文件,可以生成多个路径模式,批量传入parquet方法:
target_dates = ["2024-02-19", "2024-02-20"] # 生成每个日期对应的路径模式 path_patterns = [f"s3://data-raw-dev/GoogleAds/ad_group/cron_name=*/year=2024/month=02/*{date}.parquet" for date in target_dates] df_staging = spark.read.parquet(*path_patterns) df_staging.show()
以上方法都无需遍历客户名称,PySpark会自动扫描所有符合路径模式或过滤规则的文件,并将它们合并为一个DataFrame。
内容的提问来源于stack exchange,提问作者JanF
相关产品推荐
相关产品推荐

