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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 14:39:54