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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 21:57:21