PySpark读取目录多CSV文件转换后写回原文件的实现方法
PySpark 逐CSV处理并覆盖原文件实现方案
Spark 默认的分布式写入机制会在指定输出路径下生成多个part-*开头的分片文件,本身不支持直接将处理后的数据覆盖写入单个源CSV文件。要实现修改指定列后回写原文件的需求,不要一次性读取整个Folder目录下的所有文件——这种方式会丢失数据与源文件的映射关系,处理完成后无法拆分回原文件。最稳妥的实现方式是逐文件遍历处理+临时目录中转+文件替换。
核心实现逻辑
- 遍历目标目录,过滤得到所有需要处理的CSV文件路径
- 对单个文件独立执行读取、列值修改逻辑,处理时将DataFrame重分区为1,保证仅生成1个输出分片
- 将处理结果写入当前文件专属的临时目录,避免和其他文件、原文件冲突
- 提取临时目录中生成的CSV分片,替换掉原始文件,最后清理临时目录
可直接复用的代码示例
import os import shutil from pyspark.sql import SparkSession from pyspark.sql.functions import col # 初始化Spark会话 spark = SparkSession.builder.appName("CSVInPlaceProcess").getOrCreate() # 目标CSV存放目录 target_folder = "Folder" # 自定义列处理逻辑,按需修改即可 # 示例逻辑:将num列的值统一乘以2,content列转为全小写 def modify_columns(df): return df.withColumn("num", col("num") * 2).withColumn("content", col("content").lower()) # 逐文件处理 for filename in os.listdir(target_folder): # 只处理.csv后缀文件,过滤临时文件、隐藏文件 if not filename.endswith(".csv"): continue original_path = os.path.join(target_folder, filename) # 临时输出目录,加专属后缀避免冲突 temp_dir = os.path.join(target_folder, f"{filename}_process_tmp") # 读取单文件,有表头就保留header=True,无表头改为False source_df = spark.read.csv(original_path, header=True, inferSchema=False) # 执行列修改 processed_df = modify_columns(source_df) # 重分区为1后写入临时目录,覆盖模式避免历史残留报错 processed_df.coalesce(1).write.csv(temp_dir, header=True, mode="overwrite") # 找到临时目录里的part开头的结果分片 result_part = None for tmp_file in os.listdir(temp_dir): if tmp_file.startswith("part-") and tmp_file.endswith(".csv"): result_part = os.path.join(temp_dir, tmp_file) break # 替换原文件+清理临时目录 if result_part: shutil.move(result_part, original_path) shutil.rmtree(temp_dir) spark.stop()
注意事项
coalesce(1)仅适合单个CSV文件体积较小(一般10G以内)的场景,如果单文件体积过大,强行合并为1个分区会把所有计算压力集中到单个Executor,性能会严重下降;大文件场景建议换用pandas等单进程工具处理,或先合并分片文件再替换原文件- 上述代码适配本地文件系统,如果数据存放在HDFS、S3等分布式存储上,把
shutil对应的移动、删除逻辑替换为对应存储的文件操作API即可,核心处理逻辑不变 - 不推荐一次性读取全目录后再拆分回写的方案:就算用
input_file_name()函数给每行数据打上来路标签,处理时也会触发全量数据shuffle,文件数量多、数据量大的时候效率远低于逐文件处理
内容的提问来源于stack exchange,提问作者Priya p
相关产品推荐
相关产品推荐

