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

如何高效将Snowflake超大型表导出为S3单个CSV文件?

高效合并S3上的Spark CSV分片为单个文件的方案

问题根源

直接用coalesce(1)或repartition(1)将300GB数据合并为单个分区写入S3会失败,核心原因是:

  • 单个Executor需要承载全部300GB数据(即使压缩后,解压处理时内存占用会远超预期),极易触发OOM(内存溢出)、GC超时或任务超时。
  • Databricks集群的单Executor资源(内存、CPU)无法支撑如此大规模的单分区计算。

推荐方案(按效率从高到低排序)

方案1:直接合并S3上已有的分片文件(最优)

既然已经有110个分片在S3上,无需重新通过Spark读取全量数据,直接利用Shell命令或AWS CLI合并现有文件,效率最高。

步骤(Databricks中执行):

  1. 将分片文件下载到集群本地临时目录(确保集群本地磁盘有足够空间,300GB需要对应磁盘容量):
%sh
mkdir -p /tmp/csv_merge
aws s3 cp s3://your-bucket/path/to/your/shards/ /tmp/csv_merge/ --recursive
  1. 合并文件并处理重复Header:
%sh
# 保留第一个文件的Header,其余文件从第二行开始追加
awk 'NR==1 || FNR>1' /tmp/csv_merge/*.csv > /tmp/merged.csv
  1. 将合并后的文件上传回S3:
%sh
aws s3 cp /tmp/merged.csv s3://your-bucket/path/to/final/merged.csv
  1. 清理临时文件:
%sh
rm -rf /tmp/csv_merge /tmp/merged.csv

如果集群本地磁盘不足,可以使用S3的批量操作,但本地下载合并的稳定性更高。

方案2:Spark先合并为少量大分区,再Shell合并

如果需要重新处理数据(比如过滤、转换),先将DataFrame合并为少量大分区(比如4个,每个约75GB),降低Spark的单分区压力,写入S3后再合并为单个文件。

代码示例:

# 合并为4个大分区,每个分区约75GB,适配Executor资源
df = df.coalesce(4)

# 写入临时目录,生成4个大文件
temp_path = "s3://your-bucket/temp-large-shards/"
df.write.mode("overwrite").csv(
    path=temp_path,
    compression='gzip',
    sep='\t',
    emptyValue='""',
    nullValue='""',
    header=True
)

# 用Shell命令合并这4个文件(处理Header)
%sh
mkdir -p /tmp/large_shards
aws s3 cp s3://your-bucket/temp-large-shards/ /tmp/large_shards/ --recursive
# 合并时去重Header
awk 'NR==1 || FNR>1' /tmp/large_shards/*.csv.gz > /tmp/merged.csv.gz
# 上传到到最终路径
aws s3 cp /tmp/merged.csv.gz s3://your-bucket/final/merged.csv.gz
# 清理临时文件
rm -rf /tmp/large_shards /tmp/merged.csv.gz
dbutils.fs.rm(temp_path, recurse=True)

方案3:调整Spark配置强行合并单分区(不推荐)

如果必须用Spark直接输出单个文件,可以调整Executor资源配置,但风险极高,仅适合资源充足的超大集群:

# 调整Spark配置,加大Executor资源
spark.conf.set("spark.executor.memory", "64g")
spark.conf.set("spark.executor.cores", "8")
spark.conf.set("spark.driver.memory", "32g")
spark.conf.set("spark.sql.shuffle.partitions", "1")
spark.conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")

# 合并为1个分区写入
df = df.coalesce(1)
df.write.mode("overwrite").csv(
    path='s3://your-bucket/final/merged.csv',
    compression='gzip',
    sep='\t',
    emptyValue='""',
    nullValue='""',
    header=True
)

注意:即使调整配置,仍可能因数据量过大导致任务失败,优先考虑前两种方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 16:00:30