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

Scala实现Spark DataFrame两列键值对拆分为独立列

Spark/Scala 长表转宽表实现方案

该需求为典型的**键值对长表转宽表(行转列)**场景,此前查阅explode、split相关方案属于方向偏差,这类方法适用于字符串拆分场景,当前结构化键值对行转列直接使用Spark内置pivot算子即可实现,完整实现步骤如下:

1. 构造示例测试数据

先复现问题描述中的输入DataFrame结构,代码如下:

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

val spark = SparkSession.builder().master("local[*]").appName("PivotDemo").getOrCreate()
import spark.implicits._

// 构造与示例一致的输入数据
val inputDF = Seq(
  (0, "name", "James"),
  (0, "hair_color", "black"),
  (1, "name", "George"),
  (1, "hair_color", "black"),
  (2, "name", "Jack"),
  (2, "hair_color", "white"),
  (2, "eye_color", "blue")
).toDF("id", "attr_name", "attr_value")

输入数据打印效果如下:

+---+----------+----------+
| id| attr_name|attr_value|
+---+----------+----------+
|  0|      name|     James|
|  0|hair_color|     black|
|  1|      name|    George|
|  1|hair_color|     black|
|  2|      name|      Jack|
|  2|hair_color|     white|
|  2| eye_color|      blue|
+---+----------+----------+

2. 核心转换逻辑

按id字段分组,对attr_name列做透视转换,聚合取对应attr_value的值即可,最后按要求指定输出列顺序:

val resultDF = inputDF
  .groupBy("id")
  // 若属性列固定可直接传入列名列表提升性能,示例:.pivot("attr_name", Seq("name", "hair_color", "eye_color"))
  .pivot("attr_name")
  .agg(first("attr_value", ignoreNulls = true))
  .select("id", "name", "hair_color", "eye_color")

3. 输出效果验证

执行resultDF.show()打印结果,完全匹配预期:

+---+------+----------+---------+
| id|  name|hair_color|eye_color|
+---+------+----------+---------+
|  0| James|     black|     null|
|  1|George|     black|     null|
|  2|  Jack|     white|     blue|
+---+------+----------+---------+

优化提示

  • 若业务上属性列集合固定,建议在pivot方法中显式传入列名列表,可避免Spark全量扫描数据枚举所有属性值,大幅提升作业运行性能
  • 无对应属性的位置默认填充null,若需要替换为空字符串,可在结果DataFrame上调用.na.fill("")做空值处理

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 00:36:19