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

Spark SQL递归嵌套数据扁平化:PySpark实现可行性咨询

处理未知嵌套深度的PySpark数据扁平化需求

问题背景

我有一段多层嵌套的JSON数据,嵌套深度不固定:

{ "hierarchy": { "record": { "id": 1, "record": [ { "id": 2, "record": [ { "id": 3, "record": [ { "id": 4, "record": [] }, { "id": 5, "record": [] } ] } ] }, { "id": 6, "record": [ { "id": 7 } ] } ] } }, "type": "record" }

使用PySpark读取后得到的Schema和数据如下:

读取代码

df = spark.read.option("multiLine", True).json(file.json)
df.printSchema()
df.show(100,False)

初始Schema

<class 'pyspark.sql.dataframe.DataFrame'>
root
 |-- hierarchy: struct (nullable = true)
 | |-- record: struct (nullable = true)
 | | |-- id: long (nullable = true)
 | | |-- record: array (nullable = true)
 | | | |-- element: struct (containsNull = true)
 | | | | |-- id: long (nullable = true)
 | | | | |-- record: array (nullable = true)
 | | | | | |-- element: struct (containsNull = true)
 | | | | | | |-- id: long (nullable = true)
 | | | | | | |-- record: array (nullable = true)
 | | | | | | | |-- element: struct (containsNull = true)
 | | | | | | | | |-- id: long (nullable = true)
 | | | | | | | | |-- record: array (nullable = true)
 | | | | | | | | | |-- element: string (containsNull = true)
 |-- type: string (nullable = true)

初始数据展示

+--------------------------------------------------------------------------------------------------------------------------+------+
|hierarchy                                                                                                                  |type  |
+--------------------------------------------------------------------------------------------------------------------------+------+
|[[1,WrappedArray([2,WrappedArray([3,WrappedArray([4,WrappedArray()], [5,WrappedArray()])])], [6,WrappedArray([7,null])])]]|record|
+--------------------------------------------------------------------------------------------------------------------------+------+

尝试过explode操作,但无法处理未知的嵌套深度,希望将数据扁平化,得到每条record一行,包含id和父id的结果:

record_field id parent_id
=============================
record 1 null
record 2 1
record 3 2
record 4 3
record 5 3
record 6 1
record 7 6

请问是否可以通过Spark SQL(PySpark)实现该需求?


解决方案:用Spark SQL递归CTE实现

当然可以搞定!这种未知深度的嵌套树形结构,Spark SQL的递归公共表表达式(Recursive CTE) 是绝佳方案——它能自动逐层遍历所有嵌套层级,直到把所有record都展开成独立行。

核心思路

递归CTE分为两部分:

  1. 起始查询:提取顶层的record作为递归起点,父id设为null,同时保留它的子record数组
  2. 递归查询:通过LATERAL VIEW EXPLODE展开上一层的子节点,关联父节点的id,直到没有子节点时停止

具体实现

方式1:直接执行Spark SQL

先把读取的DataFrame注册为临时表,再执行递归CTE查询:

# 将DataFrame注册为临时表
df.createOrReplaceTempView("record_data")

# 执行递归SQL
result_df = spark.sql("""
WITH RECURSIVE record_hierarchy AS (
    -- 起始节点:顶层record
    SELECT
        'record' AS record_field,
        hierarchy.record.id AS id,
        CAST(NULL AS LONG) AS parent_id,
        hierarchy.record.record AS children
    FROM record_data
    UNION ALL
    -- 递归遍历子节点
    SELECT
        'record' AS record_field,
        child.id AS id,
        parent.id AS parent_id,
        child.record AS children
    FROM record_hierarchy parent
    LATERAL VIEW EXPLODE(parent.children) exploded AS child
    WHERE child.id IS NOT NULL -- 过滤无效的空节点
)
SELECT record_field, id, parent_id
FROM record_hierarchy
ORDER BY id
""")

# 查看结果
result_df.show()

方式2:纯PySpark API实现(可选)

如果更习惯用API而非SQL,也可以通过递归函数实现,但递归CTE的写法更简洁直观。

结果验证

执行后会得到完全符合预期的输出:

+------------+---+---------+
|record_field| id|parent_id|
+------------+---+---------+
|      record|  1|     null|
|      record|  2|        1|
|      record|  3|        2|
|      record|  4|        3|
|      record|  5|        3|
|      record|  6|        1|
|      record|  7|        6|
+------------+---+---------+

关键说明

  • 递归会自动适配任意嵌套深度,不管你的JSON有多少层嵌套,都能完整展开
  • WHERE child.id IS NOT NULL是为了过滤示例中类似[7,null]的无效空节点,避免产生无意义的行
  • 最终结果按id排序,和你期望的输出顺序一致

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 09:43:51