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

PySpark中如何避免重复使用已合并过的Delta CSV文件?

避免重复合并Delta CSV文件的实用方案

方案一:用文件格式标记已处理文件(你说的转Parquet可行)

这个思路就是把处理完的CSV转成Parquet,后续读取时只抓CSV文件,自然就不会重复处理了。

  • 具体操作:
    1. 读取时只加载未处理的CSV:
      # 只递归读取deltas目录下的所有CSV文件,自动排除Parquet
      dfTable = spark.read.format("csv")\
          .option("recursiveFileLookup","true")\
          .option("header", "true")\
          .load("/home/user/deltas/**/*.csv")
      
    2. 完成merge后,把本次处理的CSV转成Parquet(可以覆盖原路径或者存到其他目录),甚至直接删掉原CSV:
      # 获取本次处理的所有CSV文件路径
      processed_files = dfTable.inputFiles()
      # 逐个转成Parquet
      for file_path in processed_files:
          parquet_path = file_path.replace(".csv", ".parquet")
          spark.read.format("csv").option("header", "true").load(file_path)\
              .write.format("parquet").mode("overwrite").save(parquet_path)
          # 可选:删除原CSV,避免后续误读
          import os
          os.remove(file_path)
      
    3. 下次运行时,只会读取新增的CSV文件,已经转成Parquet的不会被加载。

方案二:维护一张处理记录表

比转格式更稳妥的方式是专门建个表记录已处理的文件,每次运行先对比筛选未处理的文件。

  • 具体操作:
    1. 先初始化一张处理记录表(第一次运行时执行):
      # 用Delta表存储处理记录,字段存文件路径和处理时间
      spark.createDataFrame([], schema="file_path string, process_time timestamp")\
          .write.format("delta").mode("overwrite").save("/home/user/processed_deltas")
      
    2. 找出所有未处理的CSV文件:
      # 获取deltas目录下所有CSV文件的路径
      all_csv_files = spark.sparkContext.wholeTextFiles("/home/user/deltas/**/*.csv").keys().collect()
      # 读取已处理的文件列表
      processed_df = spark.read.format("delta").load("/home/user/processed_deltas")
      processed_files = processed_df.select("file_path").rdd.flatMap(lambda x: x).collect()
      # 筛选出还没处理的文件
      unprocessed_files = [f for f in all_csv_files if f not in processed_files]
      
    3. 处理未处理的文件并更新记录:
      if unprocessed_files:
          # 读取未处理的CSV
          dfTable = spark.read.format("csv").option("header", "true").load(unprocessed_files)
          # 执行merge逻辑到主表(这里假设用id作为匹配键,按需修改)
          from delta.tables import DeltaTable
          main_table = DeltaTable.forPath(spark, "/home/user/main_table")
          main_table.alias("main")\
              .merge(dfTable.alias("delta"), "main.id = delta.id")\
              .whenMatchedUpdateAll()\
              .whenNotMatchedInsertAll()\
              .execute()
          # 把本次处理的文件记录到处理表
          from pyspark.sql.functions import current_timestamp
          new_processed_df = spark.createDataFrame(
              [(f, current_timestamp()) for f in unprocessed_files],
              schema="file_path string, process_time timestamp"
          )
          new_processed_df.write.format("delta").mode("append").save("/home/user/processed_deltas")
      
  • 好处:不用动原文件,处理历史清晰,排查问题方便。

方案三:用目录区分待处理/已处理文件

最直观的方式是把文件分目录放,新增的放pending,处理完的移到processed。

  • 具体操作:
    1. 调整目录结构:
      /home/user/deltas/
          ├── pending/  # 放新增的待处理CSV
          └── processed/  # 放处理完的文件(可以留CSV或转Parquet)
      
    2. 每次只读取pending目录的文件:
      dfTable = spark.read.format("csv")\
          .option("recursiveFileLookup","true")\
          .option("header", "true")\
          .load("/home/user/deltas/pending/")
      
    3. 处理完后把文件移到processed目录:
      import shutil
      import os
      processed_dir = "/home/user/deltas/processed/"
      pending_dir = "/home/user/deltas/pending/"
      # 遍历pending目录下的所有CSV,移动到processed并保持原结构
      for root, _, files in os.walk(pending_dir):
          for file in files:
              if file.endswith(".csv"):
                  src_path = os.path.join(root, file)
                  # 保持原目录层级
                  relative_path = os.path.relpath(root, pending_dir)
                  dest_dir = os.path.join(processed_dir, relative_path)
                  os.makedirs(dest_dir, exist_ok=True)
                  shutil.move(src_path, dest_dir)
      
  • 好处:一眼就能看出哪些文件处理过,不用额外维护记录,管理简单。

关于转Parquet的补充

你说的转Parquet方式完全可行,但要注意两点:

  • 如果后续不需要再读取这些历史数据,直接删除或移动原CSV比转格式更省事;如果需要复用历史数据,转Parquet能节省空间、提升读取性能。
  • 转Parquet时要确保列结构和原CSV一致,避免数据丢失。

内容的提问来源于stack exchange,提问作者nox8315

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 05:25:53