PySpark中如何避免重复使用已合并过的Delta CSV文件?
避免重复合并Delta CSV文件的实用方案
方案一:用文件格式标记已处理文件(你说的转Parquet可行)
这个思路就是把处理完的CSV转成Parquet,后续读取时只抓CSV文件,自然就不会重复处理了。
- 具体操作:
- 读取时只加载未处理的CSV:
# 只递归读取deltas目录下的所有CSV文件,自动排除Parquet dfTable = spark.read.format("csv")\ .option("recursiveFileLookup","true")\ .option("header", "true")\ .load("/home/user/deltas/**/*.csv") - 完成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) - 下次运行时,只会读取新增的CSV文件,已经转成Parquet的不会被加载。
- 读取时只加载未处理的CSV:
方案二:维护一张处理记录表
比转格式更稳妥的方式是专门建个表记录已处理的文件,每次运行先对比筛选未处理的文件。
- 具体操作:
- 先初始化一张处理记录表(第一次运行时执行):
# 用Delta表存储处理记录,字段存文件路径和处理时间 spark.createDataFrame([], schema="file_path string, process_time timestamp")\ .write.format("delta").mode("overwrite").save("/home/user/processed_deltas") - 找出所有未处理的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] - 处理未处理的文件并更新记录:
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。
- 具体操作:
- 调整目录结构:
/home/user/deltas/ ├── pending/ # 放新增的待处理CSV └── processed/ # 放处理完的文件(可以留CSV或转Parquet) - 每次只读取
pending目录的文件:dfTable = spark.read.format("csv")\ .option("recursiveFileLookup","true")\ .option("header", "true")\ .load("/home/user/deltas/pending/") - 处理完后把文件移到
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
相关产品推荐
相关产品推荐

