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

PySpark动态实现层级数据多代节点单行聚合方案

PySpark动态聚合多代层级关系为单行

核心思路

通过递归CTE遍历完整层级关系,收集每个顶层节点的所有后代,再根据最大后代数量动态生成CHILD1、CHILD2等字段,实现适配任意层级的动态聚合。

步骤与代码

1. 构造示例数据(模拟用户源DataFrame)

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

spark = SparkSession.builder.appName("HierarchyAgg").getOrCreate()

# 源数据:id, name, layer, parent, child(child为单个子节点ID)
data = [
    (1, "Top1", 0, None, 2),
    (2, "Child1", 1, 1, 3),
    (3, "GrandChild1", 2, 2, None),
    (4, "Top2", 0, None, 5),
    (5, "Child2", 1, 4, 6),
    (6, "GrandChild2", 2, 5, 7),
    (7, "GreatGrandChild", 3, 6, None)
]
df = spark.createDataFrame(data, ["id", "name", "layer", "parent", "child"])

2. 递归遍历层级关系

用CTE递归获取每个顶层节点(parent IS NULL)的所有后代,记录完整路径:

# 标记顶层节点并初始化路径
hierarchy_df = df.withColumn("path", F.array(F.col("id"))) \
    .withColumn("is_top", F.when(F.col("parent").isNull(), True).otherwise(False))

# 创建临时视图供递归CTE使用
hierarchy_df.createOrReplaceTempView("temp_hierarchy")

# 递归CTE遍历所有后代
recursive_df = spark.sql("""
    WITH RECURSIVE hierarchy AS (
        SELECT id, name, layer, parent, child, path, is_top
        FROM temp_hierarchy
        WHERE is_top = true
        UNION ALL
        SELECT h.id, h.name, h.layer, h.parent, h.child, array_append(hierarchy.path, h.id), false
        FROM hierarchy
        JOIN temp_hierarchy h ON hierarchy.child = h.parent
    )
    SELECT * FROM hierarchy
""")

3. 聚合后代列表并动态生成列

  • 按顶层节点分组,收集所有后代名称并排序
  • 计算最大后代数量,动态生成对应CHILDn字段:
# 分组收集每个顶层节点的后代名称(排除顶层自身)
grouped_df = recursive_df.groupBy(F.col("path")[0].alias("top_id"), F.first("name").alias("top_name")) \
    .agg(F.collect_list(F.when(F.col("layer") > 0, F.col("name"))).alias("children_list")) \
    .withColumn("children_list", F.filter(F.col("children_list"), lambda x: x.isNotNull()))

# 获取最大后代数量,动态生成CHILD列
max_child_count = grouped_df.select(F.size(F.col("children_list"))).agg(F.max("size")).first()[0]
child_columns = [F.col("children_list")[i].alias(f"CHILD{i+1}") for i in range(max_child_count)]

# 生成最终单行结果
final_df = grouped_df.select("top_id", "top_name", *child_columns)
final_df.show()

4. 处理多子节点场景(可选)

如果源数据中一个父节点有多个子节点,先聚合子节点为数组再递归:

# 聚合同一父节点的多个子节点为数组
multi_child_df = df.groupBy("id", "name", "layer", "parent") \
    .agg(F.collect_list("child").alias("children"))

multi_child_df.createOrReplaceTempView("temp_multi_hierarchy")

# 调整递归CTE适配多子节点
recursive_multi_df = spark.sql("""
    WITH RECURSIVE hierarchy AS (
        SELECT id, name, layer, parent, children, array(id) AS path, true AS is_top
        FROM temp_multi_hierarchy
        WHERE parent IS NULL
        UNION ALL
        SELECT h.id, h.name, h.layer, h.parent, h.children, array_append(hierarchy.path, h.id), false
        FROM hierarchy
        JOIN temp_multi_hierarchy h ON array_contains(hierarchy.children, h.id)
    )
    SELECT * FROM hierarchy
""")

# 后续聚合逻辑同步骤3

输出示例

+------+---------+-------+-------------+-------------------+
|top_id| top_name| CHILD1|      CHILD2|            CHILD3|
+------+---------+-------+-------------+-------------------+
|     1|     Top1|Child1|GrandChild1|               null|
|     4|     Top2|Child2|GrandChild2|GreatGrandChild|
+------+---------+-------+-------------+-------------------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 00:53:15