PySpark合并数据湖分层Parquet文件添加日期列去重实现方法
实现方案
路径规则与需求说明
- 源数据路径格式:
SYSTEM/Data/Year/Month/Date/Hist/hist.parquet,每个日期层级的Hist目录下存放对应日期的hist.parquet文件 - 处理目标:将所有日期的parquet文件合并为单个parquet文件写入目标存储路径
- 处理规则:
- 读取每个文件时新增
upload_date列,取值为文件对应路径的年/月/日,格式为YYYY-MM-DD - 所有源文件列名完全一致,读取单文件后先做行去重再参与合并
- 读取每个文件时新增
- 源表固定字段:
ID、Team、Flag、Action、Status - 已有代码框架:
for year_folder in year_folders: month_folders = dbutils.fs.ls(year_folder.path)
基于现有遍历框架的完整代码
代码适配Databricks环境(使用dbutils工具),补全了三层文件夹遍历、异常跳过、日期格式处理、去重、合并写入逻辑:
from pyspark.sql import functions as F # 配置项,根据实际路径修改 SOURCE_ROOT_PATH = "SYSTEM/Data/" TARGET_OUTPUT_PATH = "YOUR_TARGET_STORAGE_PATH/merged_result.parquet" # 初始化空结果集 merged_df = None # 获取所有年份文件夹,过滤非数字命名的脏目录 year_folders = [folder for folder in dbutils.fs.ls(SOURCE_ROOT_PATH) if folder.name.rstrip("/").isdigit()] for year_folder in year_folders: month_folders = dbutils.fs.ls(year_folder.path) for month_folder in month_folders: # 过滤非数字命名的月份目录 month_val = month_folder.name.rstrip("/") if not month_val.isdigit(): continue date_folders = dbutils.fs.ls(month_folder.path) for date_folder in date_folders: # 过滤非数字命名的日期目录 date_val = date_folder.name.rstrip("/") if not date_val.isdigit(): continue # 拼接当前parquet文件完整路径 current_file_path = f"{date_folder.path}Hist/hist.parquet" # 校验文件是否存在,跳过缺失文件避免任务报错 try: dbutils.fs.ls(current_file_path) except Exception as e: print(f"跳过不存在的文件路径: {current_file_path}") continue # 读取当前文件 current_df = spark.read.parquet(current_file_path) # 单文件内去重 current_df = current_df.dropDuplicates() # 拼接符合格式要求的upload_date,月份日期自动补前导零 year_val = year_folder.name.rstrip("/") upload_date = f"{year_val}-{month_val.zfill(2)}-{date_val.zfill(2)}" current_df = current_df.withColumn("upload_date", F.lit(upload_date).cast("date")) # 合并到总结果集 if merged_df is None: merged_df = current_df else: merged_df = merged_df.unionByName(current_df) # 可选:如果需要跨文件全局去重,打开下方注释,指定业务主键字段即可 # merged_df = merged_df.dropDuplicates(subset=["ID", "Team", "Flag", "Action", "Status"]) # 写入目标路径,coalesce(1)用于输出单个parquet文件,数据量大于100G时不建议使用,易触发OOM merged_df.coalesce(1).write.mode("overwrite").parquet(TARGET_OUTPUT_PATH)
更简洁的无循环实现
如果不需要严格遵循手写遍历的逻辑,可以直接用通配符读取所有文件,通过Spark内置函数解析路径提取日期,代码量更少,执行效率更高:
from pyspark.sql import functions as F SOURCE_PATH = "SYSTEM/Data/*/*/*/Hist/hist.parquet" TARGET_PATH = "YOUR_TARGET_STORAGE_PATH/merged_result.parquet" # 通配符一次性读取所有符合规则的parquet文件 raw_df = spark.read.parquet(SOURCE_PATH) # 新增文件路径字段用于解析日期 raw_df = raw_df.withColumn("_file_path", F.input_file_name()) # 拆分路径提取年、月、日,按路径层级倒序取对应位置值 raw_df = raw_df.withColumn("_path_arr", F.split(F.col("_file_path"), "/")) \ .withColumn("year", F.col("_path_arr").getItem(-5)) \ .withColumn("month", F.lpad(F.col("_path_arr").getItem(-4), 2, "0")) \ .withColumn("day", F.lpad(F.col("_path_arr").getItem(-3), 2, "0")) \ .withColumn("upload_date", F.concat_ws("-", F.col("year"), F.col("month"), F.col("day")).cast("date")) # 删除临时字段+去重 final_df = raw_df.drop("_file_path", "_path_arr", "year", "month", "day").dropDuplicates() # 写入目标 final_df.coalesce(1).write.mode("overwrite").parquet(TARGET_PATH)
注意事项
- 月份、日期做补零处理是为了严格符合
YYYY-MM-DD格式要求,避免1月被识别为2024-1-05这类不规范日期 - 代码中加入了目录过滤、文件存在校验逻辑,可自动跳过路径下的临时文件、隐藏文件、缺失文件,降低任务报错概率
coalesce(1)会将所有数据汇总到单个计算节点生成单文件,仅适合数据量较小的场景;如果总数据量较大,建议去掉该参数,直接写入即可,输出的多个parquet文件逻辑上仍是一个完整的数据集- 如果业务上存在跨文件重复数据,可在写入前增加全局去重逻辑,指定业务主键字段即可
内容的提问来源于stack exchange,提问作者Scope
相关产品推荐
相关产品推荐

