如何将超大Spark DataFrame写入AWS S3单CSV文件?
超大Spark DataFrame写入单个S3 CSV文件的可行方案
方案1:Spark并行写临时文件 + Boto3多线程合并
这是最通用、无集群依赖的方案,核心思路是先让Spark并行写入多份小文件,再通过AWS SDK合并成单个文件,彻底避免单节点内存过载:
- 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/")
- 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} )
- 清理临时文件
合并完成后删除临时目录,避免占用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
相关产品推荐
相关产品推荐

