Azure Databricks无递归CTE下用Spark/Python/SQL查子节点根父ID
层级根节点查询方案(Azure Databricks 适配)
场景说明
在Azure Databricks环境中,需从存储父子层级关联关系的源表中,查询每个子节点对应的最顶层根父ID。
环境限制:当前环境不支持CONNECT BY层级查询语法,也不支持递归CTE语法。
相关参考示例:
- 源表结构:

- 对应树形结构:

- 预期输出结果:

实现思路
由于无法使用原生递归查询语法,采用迭代追溯的方式实现:
- 初始状态下,每个节点的当前关联上级为直接父ID
- 每一轮迭代判断当前关联的上级是否为顶层根(判定规则:该ID不存在于子ID列中,即没有更上层的父节点)
- 对未找到顶层根的节点,再往上追溯一层父ID,重复步骤2直到所有节点都匹配到顶层根
方案1:PySpark实现(无第三方依赖)
代码兼容所有Databricks Runtime版本,执行逻辑透明可控,替换源表名即可直接使用:
from pyspark.sql import functions as F from pyspark.sql.types import BooleanType # 替换为实际业务中的源表名 source_df = spark.table("hierarchy_table").select( F.col("child_id").alias("node_id"), F.col("parent_id") ) # 广播所有子节点ID,用于快速判断是否到达顶层根 all_child_set = set(source_df.select("node_id").distinct().rdd.map(lambda row: row[0]).collect()) bc_child_set = spark.sparkContext.broadcast(all_child_set) @F.udf(returnType=BooleanType()) def is_top_root(pid): return pid not in bc_child_set.value # 初始化映射关系:每个节点初始关联直接父ID current_mapping = source_df.withColumn("root_id", F.col("parent_id")) while True: # 拆分已找到根、未找到根的两部分数据 rooted_part = current_mapping.filter(is_top_root(F.col("root_id"))) unrooted_part = current_mapping.filter(~is_top_root(F.col("root_id"))) # 所有节点都找到根时终止迭代 if unrooted_part.count() == 0: break # 未找到根的节点向上追溯一层 parent_ref = source_df.select( F.col("node_id").alias("root_id"), F.col("parent_id").alias("next_root_id") ) updated_unrooted = unrooted_part.join(parent_ref, on="root_id", how="left")\ .drop("root_id")\ .withColumnRenamed("next_root_id", "root_id") # 合并数据进入下一轮迭代 current_mapping = rooted_part.unionByName(updated_unrooted) # 输出最终结果,结构和预期完全一致 final_result = current_mapping.select("node_id", "root_id").orderBy("node_id") final_result.show()
方案2:纯Spark SQL实现(无UDF)
如果偏好纯SQL实现,可通过迭代临时视图的方式完成,逻辑和PySpark方案完全一致:
-- 初始化:缓存所有子节点ID CREATE OR REPLACE TEMP VIEW all_child_nodes AS SELECT DISTINCT child_id AS node_id FROM hierarchy_table; -- 初始化节点-上级映射表 CREATE OR REPLACE TEMP VIEW node_mapping AS SELECT child_id AS node_id, parent_id AS current_root_id FROM hierarchy_table; -- 迭代开关 SET unrooted_exists = 1; WHILE ${unrooted_exists} = 1 DO -- 提取已经找到顶层根的节点 CREATE OR REPLACE TEMP VIEW confirmed_root AS SELECT node_id, current_root_id AS root_id FROM node_mapping WHERE current_root_id NOT IN (SELECT node_id FROM all_child_nodes); -- 对未找到根的节点向上追溯一层 CREATE OR REPLACE TEMP VIEW updated_unrooted AS SELECT m.node_id, h.parent_id AS current_root_id FROM node_mapping m INNER JOIN hierarchy_table h ON m.current_root_id = h.child_id WHERE m.current_root_id IN (SELECT node_id FROM all_child_nodes); -- 检查是否还有未追溯完成的节点 SET unrooted_exists = (SELECT COUNT(1) FROM updated_unrooted); -- 更新映射表进入下一轮 CREATE OR REPLACE TEMP VIEW node_mapping AS SELECT node_id, current_root_id FROM confirmed_root UNION ALL SELECT node_id, current_root_id FROM updated_unrooted; END WHILE; -- 查询最终结果 SELECT node_id, root_id FROM node_mapping ORDER BY node_id;
性能提示:迭代次数等于业务数据的最大层级深度,绝大多数业务场景下层级深度不会超过10层,执行效率很高。如果数据量较大,可提前对源表的
child_id、parent_id字段做缓存或建立统计信息,进一步提升关联速度。
内容的提问来源于stack exchange,提问作者Aravindhan
相关产品推荐
相关产品推荐

