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
相关产品推荐
相关产品推荐

