Azure Databricks PySpark中如何将嵌套数组JSON展开为独立行
PySpark 实现Azure Databricks遥测JSON温度数组展开方案
之前操作触发OOM和展开失败的核心原因
- 链式调用
.withColumn()时全程保留了完整Body结构体,包含未使用的wind数组、其他无关遥测字段,单条记录冗余数据占比超90%,海量4KB小文件场景下极易打满Executor内存 - 对temp数组下的多个字段分别执行
explode会生成笛卡尔积,不仅数据量无意义膨胀,还会出现测量值和传感器编号错位的问题 - 未在读取阶段做列裁剪,全量加载所有字段后再做计算,内存开销是裁剪后的数倍
可直接落地的PySpark代码
注:代码中temp结构体下的子字段名可根据自动推断出的实际Schema直接替换,无需额外修改类型
from pyspark.sql import functions as F # 第一步:读取JSON时直接做列裁剪,*从源头丢弃所有不需要的字段,是解决OOM的核心* source_df = spark.read.json( path="abfss://<your-container>@<your-storage-account>.dfs.core.windows.net/<path-to-json>" # 小文件过多时可开启递归读取,不需要可注释 # ,recursiveFileLookup=True ).select( F.col("EnqueuedTimeUtc").alias("Measurement Time"), F.col("Body.telemetry.temp").alias("temp_records") ) # 第二步:仅对整个temp数组做一次explode,避免多次explode产生笛卡尔积 exploded_df = source_df.select( "Measurement Time", F.explode("temp_records").alias("single_temp_record") ) # 第三步:提取结构体字段映射为最终要求的输出列 final_result_df = exploded_df.select( F.lit("temperature").alias("Measurement Type"), F.col("single_temp_record.bankNumber").alias("Bank Number"), # 替换为Schema中对应的测点组号字段名 F.col("single_temp_record.sensorNumber").alias("Sensor Number"), # 替换为Schema中对应的传感器编号字段名 F.col("single_temp_record.tempValue").alias("Temperature Measurement"), # 替换为Schema中对应的温度值字段名 "Measurement Time" ) # 可选优化:根据集群规模调整分区数,避免小文件过多拖慢写入速度 # final_result_df = final_result_df.repartition(100)
额外性能调优适配Azure Databricks环境
- 集群配置层面可设置
spark.sql.files.maxPartitionBytes=268435456(即256MB),合并过多小分片减少任务调度开销 - 不需要保留原始Body字段的场景下,绝对不要在中间计算过程中携带该列,结构体嵌套的数组字段内存序列化开销远高于平铺列
- 如果需要同时处理wind数组,对wind数组单独做一次explode后和温度结果做union即可,不要在同一份DataFrame里同时explode两个数组
内容的提问来源于stack exchange,提问作者4DoorToyota
相关产品推荐
相关产品推荐

