Spark DataFrame复杂JSON扁平化及产品激活状态获取技术问询
问题描述
我在将JSON数据扁平化为表格视图时遇到了问题,具体情况如下:
- 我通过以下代码将JSON文件读入Spark DataFrame:
df = spark.read.option("multiLine","true").json('path/to/file.json')
- 该JSON包含嵌套数组结构,核心难点在于处理
attribute和value元素:需要将attribute的值作为列名,value的值作为行数据;同时存在重名的attribute(如Created_Date),还有Num、Description这类直接作为列的元素。 - 我不需要全部数据,全量扁平化或按需选择元素的方案均可接受,目标输出类似如下表格:
| ID | Product Created Date | Active_Item | Parent Num |
|---|---|---|---|
| XYZ_12345 | 2021-01-10 | Y | 123 |
| ABC_54321 | 2024-01-02 | Y | 543 |
补充问题
给定以下简化后的JSON,如何获取每个产品是否激活?注意Active_Item在Product Properties数组中的索引不固定:
{ "products": [ { "Name": "Product #1", "Product Properties": [ { "attribute": "Created_Date", "value": "2024-01-10" }, { "attribute": "Active_Item", "value": "Y" } ] }, { "Name": "Product 2", "Product Properties": [ { "attribute": "Created_Date", "value": "2024-01-02" }, { "attribute": "Modified_Date", "value": "2024-01-03" }, { "attribute": "Active_Item", "value": "Y" } ] } ] }
解决方案
一、按需提取字段实现扁平化
如果只需要目标表格中的指定字段,针对性提取比全量扁平化效率更高,步骤如下:
- 展开嵌套的
Product Properties数组 - 过滤出需要的
attribute字段 - 通过
pivot将属性名转为列名 - 合并原表中的顶层字段(如ID、Parent Num)
示例代码:
from pyspark.sql import functions as F # 读取JSON数据 df = spark.read.option("multiLine","true").json('path/to/file.json') # 展开Product Properties数组,保留顶层字段 exploded_df = df.select( "ID", "Parent_Num", F.explode("Product_Properties").alias("props") ) # 提取属性名和属性值,过滤目标字段 filtered_df = exploded_df.select( "ID", "Parent_Num", F.col("props.attribute").alias("attr"), F.col("props.value").alias("val") ).filter(F.col("attr").isin("Created_Date", "Active_Item")) # 转列为表格格式,并重命名列名 final_df = filtered_df.groupBy("ID", "Parent_Num")\ .pivot("attr")\ .agg(F.first("val"))\ .withColumnRenamed("Created_Date", "Product Created Date")\ .select("ID", "Product Created Date", "Active_Item", "Parent_Num") final_df.show()
若存在重复的attribute(如同一产品有多个Created_Date),可将F.first()替换为F.collect_list()保留所有值,或根据业务需求取最大/最小日期。
二、快速获取Active_Item(针对补充问题)
无需展开整个数组,直接用Spark数组函数定位目标属性:
from pyspark.sql import functions as F df = spark.read.option("multiLine","true").json('path/to/file.json') # 筛选出Active_Item并提取其值 result_df = df.select( "Name", F.element_at( F.filter("Product Properties", lambda x: x.attribute == "Active_Item"), 1 ).value.alias("Active_Item") ) result_df.show()
F.filter()遍历数组,保留attribute为"Active_Item"的元素F.element_at(...,1)取筛选后数组的第一个元素(默认一个产品仅一个激活状态)- 最后提取该元素的
value作为结果列
若需处理同一产品存在多个Active_Item值的情况,可改用collect_list()收集所有值:
result_df = df.select( "Name", F.collect_list( F.when(F.col("Product Properties.attribute") == "Active_Item", F.col("Product Properties.value")) ).alias("Active_Items") )
内容的提问来源于stack exchange,提问作者simplicITy
相关产品推荐
相关产品推荐

