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

PySpark解析JSON:从Key提取Id并过滤statC属性的报错与解决

PySpark解析扁平化JSON:提取Id并过滤含statC的行

问题描述

我正在使用PySpark DataFrame API解析并扁平化JSON数据,需要从JSON的Key/属性中提取数据元素'Id',同时仅过滤出包含'statC'属性且statC内存在'newValue'的行。尝试用explode函数展开JSON对象时出现报错,求可行的解决方法。

输入JSON

{
  "changes": {
    "1": [
      {
        "Name": "ABC-1",
        "statC": {
          "newValue": 10
        },
        "column": {
          "notDone": true,
          "newStatus": "10071"
        }
      }
    ],
    "2": [
      {
        "Name": "ABC-2",
        "added": true
      }
    ],
    "3": [
      {
        "Name": "ABC-3",
        "column": {
          "notDone": true,
          "newStatus": "10071"
        }
      }
    ],
    "4": [
      {
        "Name": "ABC-4",
        "statC": {
          "newValue": 40
        }
      }
    ],
    "5": [
      {
        "Name": "ABC-5",
        "statC": {
          "newValue": 50
        },
        "column": {
          "notDone": false,
          "done": true,
          "newStatus": "13685"
        }
      }
    ],
    "6": [
      {
        "Name": "ABC-61",
        "added": true
      },
      {
        "Name": "ABC-62",
        "statC": {
          "oldValue": 60
        }
      }
    ],
    "7": [
      {
        "Name": "ABC-70",
        "added": true
      },
      {
        "Name": "ABC-71",
        "statC": {
          "newValue": 70
        }
      },
      {
        "Name": "ABC-72",
        "statC": {
          "newValue": 75
        }
      }
    ]
  },
  "startTime": 1666188060000,
  "endTime": 1667347140000,
  "activatedTime": 1666188126953,
  "now": 1667294686212
}

期望输出

Id  Name   statC_NewValue 
1   ABC-1  10      
4   ABC-4  40     
5   ABC-5  50 
7   ABC-71 70
7   ABC-72 75

我的代码

from pyspark.sql.functions import * 
rawDF = spark.read.json([f"abfss://{pADLSContainer}@{pADLSGen2}.dfs.core.windows.net/{pADLSDirectory}/InputFile.json"], multiLine = "true")

idDF = rawDF.select(explode("changes").alias("changes_json"))

报错信息

AnalysisException: cannot resolve 'explode(changes)' due to data type mismatch: input to function explode should be array or map type, not struct.


解决方法

报错核心原因是changes字段为Struct类型,而explode仅支持Array或Map类型输入。需先将Struct转换为键值对形式,再逐步展开过滤,具体实现代码如下:

from pyspark.sql.functions import * 

# 读取原始JSON数据
rawDF = spark.read.json([f"abfss://{pADLSContainer}@{pADLSGen2}.dfs.core.windows.net/{pADLSDirectory}/InputFile.json"], multiLine = "true")

# 1. 将changes Struct转换为键值对Map并展开,获取Id(原Struct的键)和对应的数据数组
changes_cols = rawDF.select("changes.*").columns
idDF = rawDF.select(
    explode(
        map_from_entries(
            array(*[struct(lit(c).alias("key"), col(f"changes.{c}").alias("value")) for c in changes_cols])
        )
    ).alias("Id", "data_list")
)

# 2. 展开数据数组中的每个元素
expandedDF = idDF.select(col("Id"), explode(col("data_list")).alias("data"))

# 3. 过滤出包含statC且statC存在newValue的行,提取目标字段
resultDF = expandedDF.filter(
    col("data.statC").isNotNull() & 
    col("data.statC.newValue").isNotNull()
).select(
    col("Id"),
    col("data.Name").alias("Name"),
    col("data.statC.newValue").alias("statC_NewValue")
)

# 查看结果
resultDF.show()

代码说明

  • 步骤1:通过rawDF.select("changes.*").columns获取changes下的所有键(如1、2、7),将每个键和对应的值封装为Struct,再转换为Map后用explode展开,得到Id和对应的数据数组。
  • 步骤2:用explode展开数据数组,得到每个独立的对象条目。
  • 步骤3:过滤掉无statC或statC中无newValue的行,提取所需字段并命名,最终得到期望格式的结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 19:25:30