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

Spark DataFrame复杂JSON扁平化及产品激活状态获取技术问询

问题描述

我在将JSON数据扁平化为表格视图时遇到了问题,具体情况如下:

  1. 我通过以下代码将JSON文件读入Spark DataFrame:
df = spark.read.option("multiLine","true").json('path/to/file.json')
  1. 该JSON包含嵌套数组结构,核心难点在于处理attribute和value元素:需要将attribute的值作为列名,value的值作为行数据;同时存在重名的attribute(如Created_Date),还有Num、Description这类直接作为列的元素。
  2. 我不需要全部数据,全量扁平化或按需选择元素的方案均可接受,目标输出类似如下表格:
IDProduct Created DateActive_ItemParent Num
XYZ_123452021-01-10Y123
ABC_543212024-01-02Y543

补充问题

给定以下简化后的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"
            }
        ]
    }
]
}

解决方案

一、按需提取字段实现扁平化

如果只需要目标表格中的指定字段,针对性提取比全量扁平化效率更高,步骤如下:

  1. 展开嵌套的Product Properties数组
  2. 过滤出需要的attribute字段
  3. 通过pivot将属性名转为列名
  4. 合并原表中的顶层字段(如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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 00:50:30