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

如何将超大Spark DataFrame写入AWS S3单CSV文件?

超大Spark DataFrame写入单个S3 CSV文件的可行方案

方案1:Spark并行写临时文件 + Boto3多线程合并

这是最通用、无集群依赖的方案,核心思路是先让Spark并行写入多份小文件,再通过AWS SDK合并成单个文件,彻底避免单节点内存过载:

  1. Spark并行写入临时分区文件
    跳过coalesce(1),直接按集群最优分区数写入S3临时目录。如果需要控制分区数量,用repartition(n)调整(n根据集群CPU核数/内存设置,比如20-100不等):
# 写入带表头的临时CSV文件
csv_df.write.mode("overwrite")\
    .option("header", "true")\
    .csv("s3://your-bucket/temp-csv-staging/")
  1. Boto3多线程合并为单个文件
    利用S3的Multipart Upload特性,并行读取临时文件并合并成目标文件,全程无需加载所有数据到单节点内存:
import boto3
from concurrent.futures import ThreadPoolExecutor

s3_client = boto3.client("s3")
BUCKET = "your-bucket"
TEMP_PREFIX = "temp-csv-staging/"
TARGET_KEY = "final_single_file.csv"

# 筛选出所有有效的CSV分区文件(排除_SUCCESS、.crc等无关文件)
resp = s3_client.list_objects_v2(Bucket=BUCKET, Prefix=TEMP_PREFIX)
part_keys = [obj["Key"] for obj in resp["Contents"] if obj["Key"].endswith(".csv")]

# 初始化Multipart Upload
upload_resp = s3_client.create_multipart_upload(Bucket=BUCKET, Key=TARGET_KEY)
upload_id = upload_resp["UploadId"]
upload_parts = []

# 多线程上传每个分区文件作为Multipart的一部分
def upload_part(part_idx, part_key):
    with s3_client.get_object(Bucket=BUCKET, Key=part_key) as part_obj:
        part_resp = s3_client.upload_part(
            Bucket=BUCKET,
            Key=TARGET_KEY,
            UploadId=upload_id,
            PartNumber=part_idx + 1,
            Body=part_obj["Body"]
        )
        upload_parts.append({"PartNumber": part_idx + 1, "ETag": part_resp["ETag"]})

with ThreadPoolExecutor(max_workers=10) as executor:
    executor.map(upload_part, range(len(part_keys)), part_keys)

# 完成合并,生成单个文件
s3_client.complete_multipart_upload(
    Bucket=BUCKET,
    Key=TARGET_KEY,
    UploadId=upload_id,
    MultipartUpload={"Parts": upload_parts}
)
  1. 清理临时文件
    合并完成后删除临时目录,避免占用S3存储空间:
def delete_temp_obj(obj_key):
    s3_client.delete_object(Bucket=BUCKET, Key=obj_key)

with ThreadPoolExecutor(max_workers=10) as executor:
    executor.map(delete_temp_obj, part_keys + [f"{TEMP_PREFIX}_SUCCESS"])

方案2:EMR集群用Hadoop getmerge命令(仅限EMR环境)

如果你的Spark运行在AWS EMR上,可以直接用Hadoop内置的getmerge命令,它会分布式合并S3上的分区文件,无需手动处理多线程:

import subprocess

# 在EMR集群执行getmerge命令
subprocess.run([
    "hadoop", "fs", "-getmerge",
    "s3://your-bucket/temp-csv-staging/",
    "s3://your-bucket/final_single_file.csv"
], check=True)

这个命令会自动处理文件合并,且是分布式执行,不会给单个节点带来内存压力。

方案3:自定义Spark输出格式(进阶)

可以编写自定义Hadoop OutputFormat,让Spark并行写入同一个S3文件,但需要处理并发写入的冲突,实现复杂度较高,仅适合有Hadoop开发经验的用户,这里不展开。

核心注意事项

  • 绝对避免coalesce(1):强制合并到单个分区会把全量数据压到单节点,必然触发内存不足。
  • S3无传统追加模式,但Multipart Upload等效:S3标准对象不支持追加写入,但Multipart Upload可以将多个文件合并成单个对象,效果一致且更可靠。
  • 分区数要合理:Spark写临时文件时,分区数不宜过多(避免小文件泛滥)也不宜过少(降低并行度),建议根据集群核数设置为核数的2-4倍。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 09:48:20