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

如何在Spark中将处理后的数据保存为单个CSV文件至S3

解决Spark处理后保存单个CSV到S3的问题

为什么会生成多个part文件?

Spark是分布式计算框架,DataFrame的每个分区都会输出一个对应的part文件。哪怕你合并了DataFrame,只要DataFrame存在多个分区,输出时就会生成多个part文件。而append模式只是把新数据追加到目标路径,并不会合并已有分区的文件。

具体解决方案

1. 合并DataFrame分区后保存

通过coalesce(1)或repartition(1)把DataFrame的分区数缩减为1,这样输出就只会生成一个part文件:

  • coalesce(1):尽量减少分区,不会触发数据shuffle,性能更优,适合大数据量场景
  • repartition(1):强制将所有数据shuffle到一个分区,适合需要重新调整分区的情况

示例代码:

# 假设df_total是你处理后的最终DataFrame
df_total.coalesce(1) \
  .write \
  .option("header", "true") \
  .option("encoding", "UTF-8") \
  .mode("overwrite")  # 用overwrite覆盖目标路径,避免多次写入生成多个文件
  .csv("s3a://news/test_single")

2. 重命名生成的part文件(可选)

上面的方法会生成类似part-00000-xxxx.csv的文件,如果需要指定自定义文件名(比如final_result.csv),可以用Hadoop API完成重命名:

用PySpark结合Hadoop API的示例

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

# 定义输出路径
output_path = "s3a://news/test_single"

# 保存合并后的DataFrame
df_total.coalesce(1).write.option("header", "true").mode("overwrite").csv(output_path)

# 获取Hadoop文件系统实例
fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration())
path_obj = spark._jvm.org.apache.hadoop.fs.Path(output_path)

# 筛选出part文件
part_files = []
for status in fs.listStatus(path_obj):
    file_name = status.getPath().getName()
    if "part-" in file_name and file_name.endswith(".csv"):
        part_files.append(status.getPath().toString())

# 重命名第一个part文件为目标文件名
if part_files:
    target_file = f"{output_path}/final_result.csv"
    fs.rename(spark._jvm.org.apache.hadoop.fs.Path(part_files[0]), spark._jvm.org.apache.hadoop.fs.Path(target_file))
    # 可选:删除原part文件和_SUCCESS标记文件
    fs.delete(spark._jvm.org.apache.hadoop.fs.Path(part_files[0]), False)
    fs.delete(spark._jvm.org.apache.hadoop.fs.Path(f"{output_path}/_SUCCESS"), False)

注意事项

  • 性能风险:把所有数据集中到一个分区,会导致单个Executor节点承载全部数据,数据量过大时可能出现内存溢出或处理缓慢,这种方式更适合中小数据集。超大数据量的话,建议不要强行合并,后续可通过S3的批量操作完成文件合并。
  • append模式的局限:如果使用append模式,即使合并分区,每次写入都会新增一个part文件,所以要得到单个文件,建议用overwrite模式一次性写入。
  • 权限问题:确保Spark作业拥有S3路径的读写权限,包括文件重命名的操作权限。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 09:32:59