Databricks读取JSON生成DataFrame仅单分区问题咨询
问题产生原因
- 核心原因是读取JSON时开启了
multiline=True:Spark在multiline模式下会将整个文件判定为单个完整JSON实体,不会按文件块做输入分片,因此读入阶段直接生成1个分区。你之前调整的shuffle partitions参数、AQE开关、分区配置都只作用于shuffle阶段,而文件读入是窄依赖流程,不触发shuffle,这些参数完全不生效。 - 后续的JSON解析、扁平化操作(包括
explode、自定义parse_json)默认都是窄依赖,不会自动触发数据重分布,所有7.5亿条记录的计算全部压在单个executor上,必然出现执行耗时过长、OOM报错。 - 你的源JSON是顶层包含多层嵌套数组的单对象结构,multiline模式下Spark无法自动识别数组内元素边界做内部切分,只能全量加载到单分区处理。
有效解决方案
按落地成本从低到高排序:
- 方案1:读取后立即显式重分区(最快落地,适配你当前的代码逻辑)
保留现有读取逻辑不变,读完数据后第一时间调用repartition方法将数据打散到多个分区,再执行后续扁平化操作。分区数按集群总CPU核数的23倍设置即可,比如集群总核数为128时,设置256384个分区。示例代码:
由于你的原始文件仅500行,体积很小,重分区的shuffle开销可以忽略不计。# 原有读取逻辑 df = spark.read.load(Filename, format='json', multiline=True, encoding = 'UTF-8', schema= schema) # 立即重分区打散数据,不要用coalesce,coalesce无法实现单分区到多分区的全量打散 df = df.repartition(256) # 后续再执行explode、parse_json等扁平化操作 - 方案2:预处理源文件为JSON Lines格式,从读入阶段实现多分区
如果业务允许调整源文件格式,可以提前把源文件处理为JSON Lines格式(每行对应一个独立的JSON对象,比如把顶层abc数组下的每个元素拆为单独一行),读取时去掉multiline=True参数,Spark会自动按文件块切分输入分片,读入阶段直接生成多分区,不需要额外重分区。 - 方案3:逐层打散避免数据倾斜(适合重分区后仍有单task负载过高的场景)
如果重分区后做多层explode时还是出现单task数据倾斜,可以在第一次explode顶层数组时添加随机盐值,按盐值分区后再做后续嵌套层的解析,示例逻辑:from pyspark.sql import functions as F df = spark.read.load(Filename, format='json', multiline=True, encoding = 'UTF-8', schema= schema) # 炸开第一层abc数组,同时生成0-255范围的随机盐 df = df.select( F.explode("abc").alias("abc_data"), F.floor(F.rand() * 256).alias("salt") ) # 按盐值重分区,保证数据均匀分布到所有task df = df.repartition(256, "salt") # 后续再逐层炸开b、c等嵌套数组,完成扁平化逻辑
内容的提问来源于stack exchange,提问作者rdobbss
相关产品推荐
相关产品推荐

