如何在Scala版Spark DataFrame中将JSON对象的键作为行实体
解决方法
方法一:使用Spark内置函数(推荐)
Spark 3.0+提供的map_entries函数可直接将Struct类型转换为键值对数组,结合explode展开后即可提取所需字段:
import org.apache.spark.sql.functions._ val jsonString = """{"key1": "value1", "key2": {"nested1": "value2", "nested2": {"deeplyNested": "value3"}}}""" val df = spark.read.json(Seq(jsonString).toDS) // 将key2转换为键值对数组,展开后提取目标字段 val resultDF = df .select( $"key1", explode(map_entries($"key2")).alias("kv") // 把key2的结构拆成(key, value)数组并展开 ) .select( $"key1", $"kv.key".alias("PARAMETER"), $"kv.value.deeplyNested".alias("deeplyNested") // 从嵌套value中提取深层字段 ) .filter($"deeplyNested".isNotNull) // 过滤出存在deeplyNested的行 resultDF.show()
执行后输出:
+------+----------+-------------+ | key1| PARAMETER|deeplyNested| +------+----------+-------------+ |value1| nested2| value3| +------+----------+-------------+
方法二:使用json4s自定义UDF(适用于复杂嵌套场景)
如果嵌套结构更复杂,或需要更灵活的JSON解析,可结合json4s编写自定义UDF:
首先确保项目依赖中包含json4s(sbt项目添加:libraryDependencies += "org.json4s" %% "json4s-native" % "4.0.6")
import org.apache.spark.sql.functions._ import org.json4s._ import org.json4s.native.JsonMethods._ // 自定义UDF:解析key2的JSON字符串,提取含deeplyNested的键值对 val extractNestedUdf = udf((key2: String) => { implicit val formats = DefaultFormats val json = parse(key2) json.extract[Map[String, JValue]].flatMap { case (k, v) => (v \ "deeplyNested").extractOpt[String].map(v => (k, v)) }.toList }) val jsonString = """{"key1": "value1", "key2": {"nested1": "value2", "nested2": {"deeplyNested": "value3"}}}""" val df = spark.read.json(Seq(jsonString).toDS) .select($"key1", to_json($"key2").alias("key2_str")) // 将key2转为JSON字符串 val resultDF = df .select( $"key1", explode(extractNestedUdf($"key2_str")).alias("kv") ) .select( $"key1", $"kv._1".alias("PARAMETER"), $"kv._2".alias("deeplyNested") ) resultDF.show()
关键说明
- 方法一中的
map_entries是Spark 3.x新增API,无需手动枚举字段名,直接将Struct转为键值对数组; explode用于将数组行拆分为多行,实现逐个处理每个键值对;- 最终通过
filter筛选出有效行,得到目标结构的DataFrame。
内容的提问来源于stack exchange,提问作者bbgghh
相关产品推荐
相关产品推荐

