如何在DataFrame中处理嵌套JSON并提取指定属性后合并回原表
问题描述
给定嵌套JSON结构:
{ "positions": { "node": "abc" }, "submissions" :{ "submissionOffsets":[ { "attributeName": "sample1", "attributeValue": 1224 }, { "attributeName": "sample2", "attributeValue": 1224 }, { "attributeName": "sample3", "attributeValue": 1224 }, { "attributeName": "sample4", "attributeValue": 1224 } ] } }
需读取其中的submissionOffsets数组,根据指定的attributeName(例如sample1)提取对应的attributeName和attributeValue,输出需满足以下结构:
匹配场景
{ "positions": { "node": "abc" }, "submissions" :{ "submissionOffsets":[ { "attributeName": "sample1", "attributeValue": 1224 }, { "attributeName": "sample2", "attributeValue": 1224 }, { "attributeName": "sample3", "attributeValue": 1224 }, { "attributeName": "sample4", "attributeValue": 1224 } ] }, "attributeName": "sample1", "attributeValue": 1224 }
不匹配场景
{ "positions": { "node": "abc" }, "submissions" :{ "submissionOffsets":[ { "attributeName": "sample1", "attributeValue": 1224 }, { "attributeName": "sample2", "attributeValue": 1224 }, { "attributeName": "sample3", "attributeValue": 1224 }, { "attributeName": "sample4", "attributeValue": 1224 } ] }, "attributeName": null, "attributeValue": 0.00 }
要求用DataFrame实现,尝试过对submissions.submissionOffsets执行explode操作,筛选指定属性名和值,但仅得到单列数据,无法合并回原DataFrame。
解决方案
核心逻辑:直接从嵌套数组中筛选目标元素,提取字段后作为新列合并回原DataFrame,保留原始嵌套结构。以下以Spark DataFrame为例(Pandas逻辑可参考):
1. 加载JSON数据到DataFrame
from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, lit, array_contains, element_at, filter, struct spark = SparkSession.builder.appName("ExtractTargetAttribute").getOrCreate() # 加载示例JSON json_data = """ { "positions": { "node": "abc" }, "submissions" :{ "submissionOffsets":[ { "attributeName": "sample1", "attributeValue": 1224 }, { "attributeName": "sample2", "attributeValue": 1224 }, { "attributeName": "sample3", "attributeValue": 1224 }, { "attributeName": "sample4", "attributeValue": 1224 } ] } } """ df = spark.read.json(spark.sparkContext.parallelize([json_data]))
2. 提取目标属性并合并回原数据
target_attr_name = "sample1" # 筛选匹配的属性元素,无匹配时返回指定默认值 df_with_matched = df.withColumn( "matched_attr", when( # 先判断数组中是否存在目标属性 array_contains(col("submissions.submissionOffsets.attributeName"), target_attr_name), # 从数组中筛选出匹配的第一个元素 element_at( filter(col("submissions.submissionOffsets"), lambda x: x.attributeName == target_attr_name), 1 ) ).otherwise( # 不匹配时返回指定结构 struct(lit(None).alias("attributeName"), lit(0.00).alias("attributeValue")) ) ) # 将匹配结果拆分为独立列,保留原数据结构 result_df = df_with_matched.select( col("positions"), col("submissions"), col("matched_attr.attributeName"), col("matched_attr.attributeValue") )
3. 输出结果
执行以下代码即可得到符合要求的JSON格式:
print(result_df.toJSON().collect()[0])
关键说明
- 避免使用
explode:直接用filter和element_at操作数组,不会改变原DataFrame的行结构,无需后续合并操作; - 空值处理:通过
when/otherwise明确指定不匹配场景的返回值; - 高效性:无需展开整个数组,减少数据处理量,适合大数据场景。
内容的提问来源于stack exchange,提问作者ravish kumar
相关产品推荐
相关产品推荐

