基于PySpark读取Data Lake Gen2中最新Parquet文件的需求
实现PySpark读取Azure Data Lake Gen2指定目录下的Parquet文件
需求明确
- 基础存储路径:
base_path = /mnt/mountpoint/sales/ - 目录结构:按日期层级划分,格式为
sales/yyyy/mm/dd/hh_mm_ss/ - 核心要求:
- 优先读取当日
dd目录下最新时间戳的hh_mm_ss子目录中的Parquet文件 - 必须排除文件名以
_开头的日志类文件 - 若当日
dd目录不存在,自动定位到最近存在的dd目录读取对应文件
- 优先读取当日
优化后的实现代码
from datetime import datetime, timedelta from typing import List def get_latest_valid_parquet_paths(base_path: str) -> List[str]: # 初始化当前日期为当日 current_date = datetime.now() date_dir_format = "%Y/%m/%d" # 循环查找最近存在的日期目录(yyyy/mm/dd) while True: target_date_path = f"{base_path.rstrip('/')}/{current_date.strftime(date_dir_format)}" try: # 尝试列出目录,能成功执行则说明目录存在 dbutils.fs.ls(target_date_path) break except Exception: # 目录不存在,日期往前推一天继续找 current_date -= timedelta(days=1) # 获取该日期目录下所有时间戳子目录 timestamp_dirs = dbutils.fs.ls(target_date_path) # 筛选出符合hh_mm_ss格式的有效子目录(确保是目录而非文件) valid_timestamp_dirs = [ dir_info for dir_info in timestamp_dirs if dir_info.size == 0 and dir_info.name[-1] == '/' and len(dir_info.name.strip('/')) == 8 ] # 按时间戳从新到旧排序,取最新的目录 valid_timestamp_dirs.sort( key=lambda x: datetime.strptime(x.name.strip('/'), "%H_%M_%S"), reverse=True ) latest_timestamp_dir = valid_timestamp_dirs[0].path # 获取最新目录下的所有文件,过滤出符合要求的Parquet文件 all_files = dbutils.fs.ls(latest_timestamp_dir) valid_parquet_paths = [ file.path for file in all_files if file.name.endswith('.parquet') and not file.name.startswith('_') ] return valid_parquet_paths # 使用示例:获取路径并读取Parquet文件 target_parquets = get_latest_valid_parquet_paths("/mnt/mountpoint/sales/") df = spark.read.parquet(*target_parquets)
代码逻辑解析
- 日期目录定位:从当日开始检查目标日期目录是否存在,不存在则逐天回退,直到找到第一个存在的目录
- 时间戳目录筛选:在找到的日期目录下,只保留符合
hh_mm_ss格式的子目录,再按时间戳倒序排序取最新的一个 - 有效文件过滤:在最新时间戳目录中,只保留以
.parquet结尾且不以_开头的文件,排除日志类无效文件 - 文件读取:直接用Spark读取筛选后的文件路径,无需递归遍历所有目录,提升效率
对比原函数的优化点
- 去掉冗余递归,改为定向查找目标目录,大幅减少IO操作,提升性能
- 新增日期回退逻辑,自动处理当日目录不存在的场景
- 明确过滤无效文件,避免读取日志类文件
- 逻辑更清晰,代码可读性更高
内容的提问来源于stack exchange,提问作者asd
相关产品推荐
相关产品推荐

