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

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()

代码说明

  1. 初始层:直接筛选出所有有父节点的记录,标记为parent_1,层级为1
  2. 递归层:将当前层级的父节点与原表关联,获取上一级父节点,层级自动+1,生成对应的relationship_type
  3. 终止条件:当父节点为NULL时停止递归,避免无限循环
  4. 结果排序:按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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 05:01:35