如何使用PySpark获取S3存储桶中最新文件路径内的日期值
PySpark 从S3存储桶提取最新日期文件实现方案
完整实现代码
from pyspark.sql import SparkSession import re # 初始化SparkSession,提前确保集群已配置S3访问权限 spark = SparkSession.builder.appName("fetch_latest_s3_date").getOrCreate() sc = spark.sparkContext # 如集群未默认配置S3权限,可取消下方注释自行配置 # sc._jsc.hadoopConfiguration().set("fs.s3a.access.key", "你的访问密钥AK") # sc._jsc.hadoopConfiguration().set("fs.s3a.secret.key", "你的加密密钥SK") # sc._jsc.hadoopConfiguration().set("fs.s3a.endpoint", "S3服务对应的endpoint地址") # 定义目标S3目录前缀 s3_root_path = "s3://bucketname/folderpath/" # 调用Hadoop FileSystem API获取目录下所有Parquet文件路径 URI = sc._gateway.jvm.java.net.URI Path = sc._gateway.jvm.org.apache.hadoop.fs.Path FileSystem = sc._gateway.jvm.org.apache.hadoop.fs.FileSystem fs = FileSystem.get(URI(s3_root_path), sc._jsc.hadoopConfiguration()) file_statuses = fs.globStatus(Path(f"{s3_root_path}/*/*/*/*.parquet")) all_file_paths = [status.getPath().toString() for status in file_statuses] # 解析路径内日期,排序取最大值 date_regex = re.compile(r"(\d{4})/(\d{2})/(\d{2})") path_date_mapping = {} for file_path in all_file_paths: match_result = date_regex.search(file_path) if match_result: year, month, day = match_result.groups() date_num = int(f"{year}{month}{day}") path_date_mapping[date_num] = file_path # 得到最终结果 latest_date = max(path_date_mapping.keys()) latest_file_path = path_date_mapping[latest_date]
运行效果
执行上述代码后输出结果如下,符合要求:
- 最新路径:
s3://bucketname/folderpath/2021/10/10/file.parquet - 日期变量:
date = 20211010
注意事项
- 若你的S3路径的日期分区格式不是
年/月/日结构,调整正则表达式的匹配规则即可适配 - 若目标路径下文件量级极大,可优先使用分区过滤逻辑减少遍历范围,不需要全量枚举所有文件
- 运行前请确认Spark集群已加载s3a相关依赖,且拥有对应S3路径的读权限
内容的提问来源于stack exchange,提问作者Prainika
相关产品推荐
相关产品推荐

