Spark Scala实现DataFrame中key=value格式字符串列拆分为多列
Spark Scala 拆分key-value格式字符串列为多列实现方案
核心思路
先将properties列的字符串格式处理为Map类型,再按key提取对应值生成独立列。
实现代码
1. 导入依赖包
import org.apache.spark.sql.functions._
2. 处理字符串转Map类型
// 示例测试数据构造,实际使用时替换为你的DataFrame val data = Seq( (0, "A", "{_age=10, _city=A, _sal=1000}"), (1, "B", "{_age=20, _city=B, _sal=3000, tag=XYZ}"), (2, "C", "{_city=BC, tag=ABC}") ).toDF("id", "name", "properties") // 转换properties列为Map结构 val mapDF = data.withColumn("prop_map", map_from_entries( transform( // 先去除首尾大括号,再按`, `拆分得到所有键值对 split(regexp_replace($"properties", "[{}]", ""), ", "), // 每个键值对按`=`拆分,转为Map的entry结构 kv => split(kv, "=") ) ) )
3. 按key提取生成独立列
方案A:提前已知所有key(推荐,性能更高)
直接指定需要提取的key即可:
val result = mapDF.select( $"id", $"name", $"prop_map".getItem("_age").as("_age"), $"prop_map".getItem("_city").as("_city"), $"prop_map".getItem("_sal").as("_sal"), $"prop_map".getItem("tag").as("tag") ) // 查看结果 result.show()
方案B:动态获取所有key(适合key不确定的场景)
// 聚合获取所有出现过的key val allKeys = mapDF.agg(collect_set(expr("explode(map_keys(prop_map))"))).head().getSeq[String](0) // 动态生成查询列 val selectCols = Seq(col("id"), col("name")) ++ allKeys.map(k => col("prop_map").getItem(k).as(k)) val result = mapDF.select(selectCols:_*) // 查看结果 result.show()
输出效果
和预期结果一致,不存在的key对应列值为null:
+---+----+----+-----+----+----+ | id|name|_age|_city|_sal| tag| +---+----+----+-----+----+----+ | 0| A| 10| A|1000|null| | 1| B| 20| B|3000| XYZ| | 2| C|null| BC|null| ABC| +---+----+----+-----+----+----+
内容的提问来源于stack exchange,提问作者Subhadip Roy
相关产品推荐
相关产品推荐

