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
相关产品推荐
相关产品推荐

