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

如何用Scala+Spark实现分层数据上卷生成指定结构DataFrame

Scala + Spark 层级数据上卷实现方案

需求说明

输入为key-value结构的父子层级映射数据,需要展开为每个原始id的三级上级结构,无对应层级时填充null,输出字段为id、level1、level2、level3。
样例输入映射关系:

  • e1 → e2
  • e2 → e3
  • e4 → e2
  • e5 → e4
  • e6 → e4

实现步骤

1. 初始化Spark环境与构造输入数据

以下代码可直接在本地测试运行:

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._

// 初始化SparkSession
val spark = SparkSession.builder()
  .appName("HierarchyRollup")
  .master("local[*]")
  .getOrCreate()

import spark.implicits._

// 构造输入的层级映射数据
val rawDf = Seq(
  ("e1", "e2"),
  ("e2", "e3"),
  ("e4", "e2"),
  ("e5", "e4"),
  ("e6", "e4")
).toDF("key", "parent")

2. 多层自关联实现层级上卷

采用三次左外关联的方式,逻辑清晰易理解,适合初学者快速掌握:

// 关联获取一级上级level1
val level1Df = rawDf.select($"key".alias("id"), $"parent".alias("level1"))

// 关联获取二级上级level2
val level2Df = level1Df
  .join(rawDf, level1Df("level1") === rawDf("key"), "left_outer")
  .select("id", "level1", "parent")
  .withColumnRenamed("parent", "level2")

// 关联获取三级上级level3
val resultDf = level2Df
  .join(rawDf, level2Df("level2") === rawDf("key"), "left_outer")
  .select("id", "level1", "level2", "parent")
  .withColumnRenamed("parent", "level3")

// 打印结果
resultDf.show()

3. 输出结果验证

运行代码后输出结果和预期完全匹配:

+---+------+------+------+
| id|level1|level2|level3|
+---+------+------+------+
| e1|    e2|    e3|  null|
| e2|    e3|  null|  null|
| e4|    e2|    e3|  null|
| e5|    e4|    e2|    e3|
| e6|    e4|    e2|    e3|
+---+------+------+------+

补充说明

如果实际业务中层级不固定、深度超过3级,可以改用Spark SQL的递归CTE实现,调整最大递归层级参数即可适配任意深度的层级结构。

内容的提问来源于stack exchange,提问作者pokiman

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 00:15:05