在PySpark中动态扁平化父/子层级数据集的技术需求
在PySpark中扁平化父子层级数据集(无parent_id)
实现思路
通过递归CTE遍历层级关系,生成包含层级和完整路径数组的中间结果,再动态将路径数组拆分为多列,自动适配新增子节点后的列扩展需求。
1. 准备样本数据
先构建包含父子关系的DataFrame,包含你提到的新增子节点Analyst:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, size, max as spark_max, expr spark = SparkSession.builder.appName("HierarchyFlatten").getOrCreate() # 样本数据(含新增的Analyst节点) data = [ ("VP", "CEO"), ("Manager", "VP"), ("Senior", "Manager"), ("Analyst", "Senior") ] df = spark.createDataFrame(data, ["Child", "Parent"]) df.show()
2. 递归CTE构建层级与路径数组
先自动识别根节点(不在Child列中的Parent值,即CEO),再递归遍历所有子节点,记录层级和完整路径数组:
# 自动获取根节点(无父节点的顶层节点) root_nodes = df.select("Parent").exceptAll(df.select("Child")).collect() root_node = root_nodes[0][0] if root_nodes else None # 递归CTE生成层级数据 hierarchy_df = spark.sql(f""" WITH RECURSIVE hierarchy AS ( -- 根节点初始化:Level=0,路径数组仅包含自身 SELECT '{root_node}' AS node, 0 AS Level, ARRAY('{root_node}') AS path_array UNION ALL -- 递归遍历子节点:层级+1,路径数组追加当前子节点 SELECT df.Child AS node, h.Level + 1 AS Level, array_append(h.path_array, df.Child) AS path_array FROM hierarchy h JOIN df ON h.node = df.Parent ) SELECT * FROM hierarchy """) hierarchy_df.show(truncate=False)
3. 动态生成Path列
根据路径数组的最大长度,自动创建Path 1至Path N列:
# 获取路径数组的最大长度,确定需要生成的Path列数 max_path_length = hierarchy_df.select(size(col("path_array")).alias("len")).agg(spark_max("len")).collect()[0][0] # 生成Path列的表达式 path_columns = [expr(f"path_array[{i}] AS `Path {i+1}`") for i in range(max_path_length)] # 构建最终结果集 final_df = hierarchy_df.select( col("node").alias("Node"), col("Level"), *path_columns ) final_df.show(truncate=False)
最终输出示例
| Node | Level | Path 1 | Path 2 | Path 3 | Path 4 | Path 5 |
|---|---|---|---|---|---|---|
| CEO | 0 | CEO | null | null | null | null |
| VP | 1 | CEO | VP | null | null | null |
| Manager | 2 | CEO | VP | Manager | null | null |
| Senior | 3 | CEO | VP | Manager | Senior | null |
| Analyst | 4 | CEO | VP | Manager | Senior | Analyst |
内容的提问来源于stack exchange,提问作者Prateek
相关产品推荐
相关产品推荐

