You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 18:31:29