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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 20:18:14