如何在PySpark中迭代连接层级数据生成完整层级树?
PySpark 生成层级完整路径解决方案
给定层级结构的PySpark DataFrame,我们可以通过**递归CTE(Common Table Expression)**实现任意深度的层级路径生成,为每条记录生成从自身到根节点(ParentID为null)的完整层级树。
初始数据
data1 = [("A",None, 1, 'Highest'),("B","A", 2, 'Medium'),("C","B", 3, 'Lowest'),("D","B", 3, 'Lowest')] df1 = spark.createDataFrame(data=data1, schema = ['ID','ParentID','Hierarchy','HierarchyName']) df1.show(truncate=False)
解决方案:递归CTE实现
方法1:使用Spark SQL递归CTE(推荐,代码更简洁)
# 创建临时视图 df1.createOrReplaceTempView("hierarchy_table") # 执行递归查询 result_df = spark.sql(""" WITH RECURSIVE hierarchy_cte AS ( -- 锚点成员:初始化每个节点的路径为自身信息 SELECT ID, ParentID, Hierarchy, HierarchyName, ARRAY(STRUCT(ID, Hierarchy, HierarchyName)) AS hierarchy_path, ParentID IS NULL AS is_root FROM hierarchy_table UNION ALL -- 递归成员:关联父节点,将父节点信息添加到路径头部 SELECT ht.ID, ht.ParentID, ht.Hierarchy, ht.HierarchyName, ARRAY(STRUCT(hc.ID, hc.Hierarchy, hc.HierarchyName)) || hc.hierarchy_path AS hierarchy_path, ht.ParentID IS NULL AS is_root FROM hierarchy_table ht JOIN hierarchy_cte hc ON ht.ParentID = hc.ID WHERE NOT hc.is_root ) -- 筛选已到达根节点的完整路径记录 SELECT ID, ParentID, Hierarchy, HierarchyName, hierarchy_path FROM hierarchy_cte WHERE is_root """) # 查看结果 result_df.show(truncate=False)
方法2:使用DataFrame API递归遍历
from pyspark.sql import functions as F # 初始化路径字段,标记是否为根节点 hierarchy_df = df1.withColumn( "hierarchy_path", F.array(F.struct(F.col("ID"), F.col("Hierarchy"), F.col("HierarchyName"))) ).withColumn("is_root", F.col("ParentID").isNull()) recursive_df = hierarchy_df.select("ID", "ParentID", "hierarchy_path", "is_root") while True: # 关联父节点,扩展层级路径 next_level = recursive_df.join(df1, recursive_df["ParentID"] == df1["ID"], "inner") \ .withColumn( "hierarchy_path", F.concat(F.array(F.struct(df1["ID"], df1["Hierarchy"], df1["HierarchyName"])), recursive_df["hierarchy_path"]) ) \ .select( df1["ID"].alias("ID"), df1["ParentID"].alias("ParentID"), "hierarchy_path", df1["ParentID"].isNull().alias("is_root") ) # 没有非根节点需要处理时终止循环 if next_level.filter(~F.col("is_root")).count() == 0: break recursive_df = recursive_df.union(next_level) # 关联原表字段,得到最终结果 result_df = recursive_df.filter(F.col("is_root")) \ .join(df1, "ID", "inner") \ .select(df1["ID"], df1["ParentID"], df1["Hierarchy"], df1["HierarchyName"], "hierarchy_path") result_df.show(truncate=False)
结果说明
生成的hierarchy_path字段为数组类型,包含从根节点到当前节点的完整层级结构:
- 记录A的路径:
[{"ID":"A","Hierarchy":1,"HierarchyName":"Highest"}] - 记录B的路径:
[{"ID":"A","Hierarchy":1,"HierarchyName":"Highest"},{"ID":"B","Hierarchy":2,"HierarchyName":"Medium"}] - 记录C/D的路径:
[{"ID":"A","Hierarchy":1,"HierarchyName":"Highest"},{"ID":"B","Hierarchy":2,"HierarchyName":"Medium"},{"ID":"C/D","Hierarchy":3,"HierarchyName":"Lowest"}]
内容的提问来源于stack exchange,提问作者tommyhmt
相关产品推荐
相关产品推荐

