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分为两部分:
- 起始查询:提取顶层的record作为递归起点,父id设为
null,同时保留它的子record数组 - 递归查询:通过
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
相关产品推荐
相关产品推荐

