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

Databricks中PySpark读取2GB嵌套JSON失败的解决方案咨询

处理Databricks中PySpark读取大体积多层嵌套JSON的解决方案

问题核心

2GB+的多层嵌套JSON使用multiline=True读取后,多次explode操作会导致数据量指数级膨胀,加上单大文件默认仅分配1个分区,所有计算集中在单个节点,超出资源承载能力后任务挂起失败。小文件因数据量小、分区压力低可正常运行。

优化步骤

1. 拆分大文件,破解单分区瓶颈

大体积multiline JSON默认只会被分配1个分区,先拆分文件分散计算压力:

# 读取为文本文件后拆分重写
text_df = spark.read.text("/FileStore/shared_uploads/konda/large_file.json")
# 按集群资源调整分区数,比如拆成100份
text_df.repartition(100).write.mode("overwrite").text("/tmp/split_large_json")
# 读取拆分后的文件
df = spark.read.option("multiline", True).json("/tmp/split_large_json")

2. 指定Schema读取,减少内存开销

提前定义JSON结构,避免Spark全文件扫描推断Schema,大幅降低内存占用:

from pyspark.sql.types import StructType, StructField, ArrayType, StringType, DoubleType, DateType

# 匹配你的JSON结构定义Schema
custom_schema = StructType([
    StructField("in_network", ArrayType(StructType([
        StructField("negotiation_arrangement", StringType()),
        StructField("name", StringType()),
        StructField("billing_code_type", StringType()),
        StructField("billing_code_type_version", StringType()),
        StructField("billing_code", StringType()),
        StructField("description", StringType()),
        StructField("negotiated_rates", ArrayType(StructType([
            StructField("negotiated_prices", ArrayType(StructType([
                StructField("billing_class", StringType()),
                StructField("expiration_date", DateType()),
                StructField("negotiated_rate", DoubleType()),
                StructField("negotiated_type", StringType()),
                StructField("service_code", ArrayType(StringType()))
            ]))),
            StructField("provider_groups", ArrayType(StructType([
                StructField("npi", ArrayType(StringType())),
                StructField("tin", StructType([
                    StructField("type", StringType()),
                    StructField("value", StringType())
                ]))
            ])))
        ])))
    ])))
])

# 用指定Schema读取文件
df = spark.read.schema(custom_schema).option("multiline", True).json("/FileStore/shared_uploads/konda/large_file.json")

3. 分步操作+分区调整,避免数据爆炸

每次explode后立即重新分区,分散后续计算压力:

# 先展开外层数组并重新分区
df = df.withColumn('in_network', explode('in_network')).repartition(200)
# 展开下一层数组后再次调整分区
df = df.withColumn('explode_negotiated_rates', explode('in_network.negotiated_rates')).repartition(300)
# 后续字段提取和explode操作同理,每步穿插分区调整

4. 精简中间数据,减少内存占用

提取字段时直接选择所需内容,尽早丢弃无用嵌套列:

df = df.select(
    explode('in_network').alias('in_network')
).select(
    col('in_network.negotiation_arrangement'),
    col('in_network.name'),
    col('in_network.billing_code_type'),
    col('in_network.billing_code_type_version'),
    col('in_network.billing_code'),
    col('in_network.description'),
    explode('in_network.negotiated_rates').alias('negotiated_rates')
).repartition(200)
# 分层处理剩余嵌套结构

5. 写入表时的性能优化

  • 用业务字段分区写入,降低单分区数据量:
df.write.mode("overwrite").partitionBy("billing_code_type").saveAsTable("your_db.your_table")
  • 启用Delta Lake的自动优化特性(Databricks环境):
spark.conf.set("spark.databricks.delta.optimizeWrite.enabled", "true")
spark.conf.set("spark.databricks.delta.autoCompact.enabled", "true")

df.write.mode("overwrite").format("delta").saveAsTable("your_db.your_table")

验证流程

  1. 先用limit取少量数据验证逻辑,确认Schema和数据正确性;
  2. 逐步扩大数据量,观察集群资源使用情况,调整分区数和executor配置(如增加内存、核心数)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 14:50:21