PySpark构建层级表(id-父id关系):优化百万级数据循环方案
优化PySpark层级关系构建(无GraphFrame,消除循环)
问题核心
原实现通过循环遍历单个base_id、调用collect()拉取数据到Driver节点、多次单条数据join操作,在百万级数据场景下完全浪费Spark分布式计算能力,性能极差。必须改用原生分布式递归方案,避免串行循环。
最优方案:递归CTE(Spark 3.0+)
利用Spark原生支持的递归CTE(Common Table Expressions),可以高效遍历未知深度的层级关系,全程分布式执行,无需循环或单点数据拉取。
完整实现代码
from pyspark.sql.types import StructType, StructField, StringType import pyspark.sql.functions as F # 测试数据 data2 = [('479', None), ('977', '666'), ('666', '479'), ('555', '678'), ('678', '977'), ] schema = StructType([ StructField("base_id", StringType(), True), StructField("parent_id", StringType(), True) ]) base_df = spark.createDataFrame(data=data2, schema=schema) # 创建临时表供递归CTE调用 base_df.createOrReplaceTempView("base_hierarchy") # 递归CTE查询 result_df = spark.sql(""" WITH RECURSIVE full_hierarchy AS ( -- 初始层:每个节点的直接父级(parent_1) SELECT base_id, parent_id, 1 AS level, 'parent_1' AS relationship_type FROM base_hierarchy WHERE parent_id IS NOT NULL UNION ALL -- 递归层:遍历父节点的父节点,生成parent_2、parent_3... SELECT fh.base_id, bh.parent_id, fh.level + 1 AS level, CONCAT('parent_', fh.level + 1) AS relationship_type FROM full_hierarchy fh JOIN base_hierarchy bh ON fh.parent_id = bh.base_id WHERE bh.parent_id IS NOT NULL ) SELECT base_id, parent_id, relationship_type FROM full_hierarchy ORDER BY base_id, level """) # 查看结果 result_df.show()
代码说明
- 初始层:直接筛选出所有有父节点的记录,标记为
parent_1,层级为1 - 递归层:将当前层级的父节点与原表关联,获取上一级父节点,层级自动+1,生成对应的
relationship_type - 终止条件:当父节点为
NULL时停止递归,避免无限循环 - 结果排序:按
base_id和层级排序,保证每个节点的父级关系按顺序展示
性能优势
- 分布式执行:全程利用Spark集群的分布式计算能力,百万级数据可并行处理
- 无单点瓶颈:避免
collect()拉取数据到Driver,消除串行循环的性能损耗 - 自动适配深度:无需提前知道最大层级深度,递归会自动遍历到顶级节点
输出结果
与原代码输出逻辑一致,每个base_id的所有父级关系都会被生成:
+-------+---------+-----------------+ |base_id|parent_id|relationship_type| +-------+---------+-----------------+ | 555| 678| parent_1| | 555| 977| parent_2| | 555| 666| parent_3| | 555| 479| parent_4| | 666| 479| parent_1| | 678| 977| parent_1| | 678| 666| parent_2| | 678| 479| parent_3| | 977| 666| parent_1| | 977| 479| parent_2| +-------+---------+-----------------+
内容的提问来源于stack exchange,提问作者WarBoy
相关产品推荐
相关产品推荐

