Azure Databricks中350GB JSON数据高效写入方案咨询
寻求大规模数据高效写入优化方案
当前我们有一套数据处理与写入流程,针对2.17亿条原始数据处理后得到9000万条数据,写入操作已运行3天8小时仍未完成,集群规模与配置无限制,平台为Azure Databricks,寻求更快的写入方案。
当前数据处理代码
df = emp_df.join(mem_df, emp_df.col1 == mem_df.col2) # df.count() → 217,000,000 grouped_df = df.groupBy( "col1", "col2", "col3", "col4", "col5" ) df1 = grouped_df \ .agg(f.collect_list(f.struct("col3", "col7", "col2", "col8", "col10")).alias("items")) \ .orderBy("col1","col2","items.col8") df2 = df1.withColumn("newCol1", f.struct(*df1.columns)).drop(*df1.columns) df3 = df2.withColumn("newCol2", f.struct(*df2.columns)).drop(*df2.columns)
输出数据结构示例
{ "newCol2": { "newCol1": { "col1": "val1", "col2": "val2", "col3": "val3", "col4": "val4", "col5": "val5", "col6": "val6", "col7": "val7", "col8": [{ "col4_1": "val8_1", "col4_1": "val8_2", "col4_1": "val8_3", "col4_1": "val8_4", "col4_1": "val8_5", "col4_1": "val8_6" }, { "col4_1": "val8_11", "col4_1": "val8_21", "col4_1": "val8_31", "col4_1": "val8_41", "col4_1": "val8_51", "col4_1": "val8_61" }... // 最多包含10个子元素 ] } } }
df3.count() → 90,000,000
当前写入代码
df3.select("data", "newCol2.newCol1.col1", "newCol2.newCol1.col2") \ .coalesce(1) \ .write \ .partitionBy("col1","col2") \ .mode("overwrite") \ .format("json") \ .save("abc/json_data/")
关键说明
- 分区列
col1约有3.5万个不同值,col2可存在于多个col1分区 - 数据总量约350GB
优化方案
1. 移除coalesce(1),恢复并行写入能力
coalesce(1)强制将所有数据合并到单个分区,完全丧失并行写入能力,这是写入极慢的核心原因。直接移除该操作,或通过repartition设置合理分区数(按每分区1-2GB估算,350GB数据可设置200-350个分区)。
示例:
df3.select("data", "newCol2.newCol1.col1", "newCol2.newCol1.col2") \ .repartition(300) # 按需调整分区数 .write \ .partitionBy("col1","col2") \ .mode("overwrite") \ .format("json") \ .save("abc/json_data/")
2. 优化分区策略,避免小分区泛滥
当前按col1+col2分区会产生大量小分区,增加文件系统元数据开销。可调整为:
- 仅按
col1分区,减少分区总数量 - 若必须保留
col2分区,可对col1进行哈希分桶或合并(如将多个col1值归为一个父分区),降低小分区占比
3. 替换文件格式,提升写入效率
JSON是文本格式,序列化开销大、压缩率低。建议改用Parquet/ORC列式存储格式,写入速度、压缩率均远优于JSON。若业务必须输出JSON,可先写入Parquet,再批量转换为JSON,整体效率更高。
示例(改用Parquet):
df3.select("data", "newCol2.newCol1.col1", "newCol2.newCol1.col2") \ .repartition(300) .write \ .partitionBy("col1") .mode("overwrite") .format("parquet") .option("compression", "snappy") .save("abc/parquet_data/")
4. 扩容集群配置(Azure Databricks)
- 选用高性能计算节点:优先选择DSv3/Fsv2系列节点,或使用Serverless SQL仓库适配批量写入场景
- 增加节点规模:直接扩容至30-50个大规格节点,最大化并行处理能力
- 调整Spark核心配置:
- 将
spark.sql.shuffle.partitions设置为集群核心数的2-3倍 - 增大
spark.sql.files.maxPartitionBytes至256MB/512MB,减少分区数 - 开启Parquet压缩为
snappy,降低存储体积与写入时间
- 将
5. 优化前置数据处理流程
- 移除冗余嵌套:
df2和df3的嵌套struct操作完全冗余,直接基于df1构造最终结构,减少序列化开销 - 提前过滤数据:在join/groupBy前过滤无用列与行,减少数据处理量
- 避免全局排序:
orderBy("col1","col2","items.col8")是全局排序,开销极大。若业务允许,改为分区内排序或直接移除排序操作
内容的提问来源于stack exchange,提问作者pramod sahoo
相关产品推荐
相关产品推荐

