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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 14:23:11