如何使用PySpark遍历嵌套文件夹提取blob存储中最新时间戳目录
PySpark 筛选日期目录下最新时间戳子目录实现方案
实现思路
- 先通过Hadoop FileSystem API递归枚举指定根目录下的所有子目录,避免直接全量读取数据浪费集群资源
- 对枚举得到的目录路径做拆分,提取出日期目录名和时间戳目录名两个核心字段
- 按日期字段分组,每组内取时间戳最大的对应目录路径
- 收集所有符合筛选条件的目录路径,统一读取路径下的parquet文件
实现代码
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import row_number, desc from org.apache.hadoop.fs import Path # 初始化SparkSession spark = SparkSession.builder.appName("get_latest_timestamp_dir").getOrCreate() # 配置blob存储访问权限(如有需要替换为自身存储账号密钥即可) # spark.conf.set("fs.azure.account.key.Storagename.blob.core.windows.net", "你的存储账号密钥") # 定义根目录路径 root_path = "wasbs://abcd@Storagename.blob.core.windows.net/dir1/dir2/" # 获取Hadoop FileSystem实例 hadoop_conf = spark._jsc.hadoopConfiguration() path = Path(root_path) fs = path.getFileSystem(hadoop_conf) # 递归列出所有子目录,深度可根据实际目录结构调整 status_list = fs.listStatus(path, 2) dir_paths = [] for status in status_list: if status.isDirectory(): full_path = status.getPath().toString() # 拆分路径提取日期和时间戳 path_parts = full_path.strip("/").split("/") # 路径格式为 {前缀}/日期/时间戳/,倒数第二段为时间戳,倒数第三段为日期 if len(path_parts) >= 3: date_str = path_parts[-3] timestamp_str = path_parts[-2] dir_paths.append((date_str, timestamp_str, full_path)) # 转为Spark DataFrame做分组取最大值操作 dir_df = spark.createDataFrame(dir_paths, schema=["date_str", "timestamp_str", "full_path"]) # 按日期分组,筛选每组最新时间戳对应的目录 window_spec = Window.partitionBy("date_str").orderBy(desc("timestamp_str")) latest_dirs_df = dir_df.withColumn("rn", row_number().over(window_spec)).filter("rn = 1").select("full_path") # 收集所有最新目录路径 latest_dir_list = [row.full_path for row in latest_dirs_df.collect()] # 读取所有最新目录下的parquet文件 df = spark.read.parquet(*latest_dir_list) # 后续可自行对读取到的df做业务处理 df.show()
注意事项
- 如果实际目录嵌套层级和示例不同,需要对应调整
path_parts的索引取值逻辑 - 递归列出目录时的深度参数可以根据自身目录结构调整,避免枚举不必要的层级提升执行效率
- 时间戳格式如果不是可直接字符串排序的格式,需要先转为datetime类型再做排序比较
内容的提问来源于stack exchange,提问作者Venkatesh
相关产品推荐
相关产品推荐

