如何用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
相关产品推荐
相关产品推荐

