Spark Scala迭代处理层级JSON生成层级表的技术实现问题
Spark Scala 层级数据转换解决方案
针对多列层级数据(3-5级)转成id/parent/name结构的需求,以下是无需动态命名DataFrame、自动适配层级数量的实现方案:
核心思路
- 自动识别所有层级列(比如以统一前缀命名的列,如
level1/level2...) - 从根层级开始,逐层生成节点并关联父节点ID
- 通过循环累加结果集的方式,避免动态创建变量,同时根据层级列数量自动终止循环
代码实现
假设输入DataFrame的层级列以level为前缀(如level1到level5),实际可根据业务调整列名匹配规则:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types.StringType // 模拟输入数据,实际替换为你的业务数据 val inputDF = spark.createDataFrame(Seq( ("Electronics", "Phones", "Smartphones", "Apple", null), ("Electronics", "Phones", "Smartphones", "Samsung", null), ("Electronics", "Laptops", "Gaming", "ASUS", null), ("Home", "Appliances", "Kitchen", "Refrigerators", "LG") )).toDF("level1", "level2", "level3", "level4", "level5") // 自动识别所有层级列并排序 val levelColumns = inputDF.columns.filter(_.startsWith("level")).sorted // 初始化根节点层级(第一层) var hierarchyDF = inputDF.select(col("level1").as("name")) .distinct() .withColumn("id", sha2(col("name"), 256)) // 用SHA256生成唯一ID,可替换为自增序列 .withColumn("parent", lit(null).cast(StringType)) // 循环处理后续层级,自动终止于最后一列 for (i <- 1 until levelColumns.length) { val currentLevelCol = levelColumns(i) val prevLevelCol = levelColumns(i - 1) // 获取当前层级与父层级的关联关系,过滤空节点 val currentNodes = inputDF .select(col(prevLevelCol).as("parent_name"), col(currentLevelCol).as("name")) .distinct() .filter(col("name").isNotNull) // 关联父节点ID,生成当前层级的父子结构 val currentHierarchy = currentNodes .join(hierarchyDF, currentNodes("parent_name") === hierarchyDF("name"), "inner") .select( sha2(col("name"), 256).as("id"), col("id").as("parent"), col("name") ) // 合并到总层级表 hierarchyDF = hierarchyDF.union(currentHierarchy) } // 最终结果:仅保留id、parent、name三列,去重 val finalHierarchyDF = hierarchyDF.select("id", "parent", "name").distinct() // 查看结果 finalHierarchyDF.show(false)
关键细节说明
- 自动适配层级数量:通过
levelColumns.length控制循环次数,不管是3级还是5级数据,都会自动处理到最后一列 - 避免动态命名DataFrame:用同一个变量
hierarchyDF累加所有层级的结果,无需创建多个动态命名的DataFrame - ID生成方式:示例用
sha2生成唯一ID,若需要更直观的自增ID,可替换为monotonically_increasing_id()或自定义序列;如果存在同名节点,可结合父节点名称生成唯一ID:sha2(concat(col("parent_name"), col("name")), 256) - 空值处理:过滤空的层级节点,避免生成无效的空条目
内容的提问来源于stack exchange,提问作者user3459079
相关产品推荐
相关产品推荐

