Spark DataFrame嵌套JSON字符串列基于Map替换键值的实现方法
Spark DataFrame 嵌套JSON键值替换解决方案
你需要替换嵌套JSON数组中的键和对应值,可通过Spark原生内置函数或自定义UDF两种方式实现,以下是可直接运行的完整代码:
前置依赖导入
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ import spark.implicits._
方案1:使用Spark原生函数(无UDF,性能更好)
利用from_json解析JSON字符串为结构化类型,结合你提到的typedLit传入Scala侧定义的映射表,转换完成后再转回JSON字符串:
// 你原有代码生成的df1基础上继续处理 val typedKMap = typedLit(kMap) val typedVMap = typedLit(vmap) // 定义嵌套数组的解析Schema val nestedSchema = ArrayType(StructType(Seq( StructField("fn", StringType), StructField("v1", StringType), StructField("v2", StringType) ))) val resultDf = df1.withColumn("nestedKey", when(col("nestedKey").isNotNull, from_json(col("nestedKey"), nestedSchema) // 遍历数组中每个元素做键值替换 .transform(element => struct( typedVMap(element.getField("fn")).alias(typedKMap("fn").cast(StringType)), element.getField("v1").alias(typedKMap("v1").cast(StringType)), element.getField("v2").alias(typedKMap("v2").cast(StringType)) )) .to_json ).otherwise(col("nestedKey"))) resultDf.show(truncate = false)
方案2:自定义UDF(灵活性更高,适合结构复杂的嵌套JSON)
如果嵌套结构不固定、层级更深,用UDF直接操作JSON对象会更简单:
- 先引入JSON处理依赖(以FastJSON为例,集群环境已导入可跳过)
<dependency> <groupId>com.alibaba</groupId> <artifactId>fastjson</artifactId> <version>1.2.83</version> </dependency>
- UDF实现代码:
import com.alibaba.fastjson.JSON import scala.collection.JavaConverters._ val replaceNestedJsonUdf = udf((jsonStr: String, keyMap: Map[String, String], valueMap: Map[String, String]) => { if (jsonStr == null) return null val jsonArr = JSON.parseArray(jsonStr) jsonArr.asScala.foreach(item => { val jsonObj = item.asInstanceOf[com.alibaba.fastjson.JSONObject] // 替换fn字段的取值 val originFn = jsonObj.getString("fn") if (valueMap.contains(originFn)) jsonObj.put("fn", valueMap(originFn)) // 替换键名 keyMap.foreach { case (oldKey, newKey) => if (jsonObj.containsKey(oldKey)) { val value = jsonObj.get(oldKey) jsonObj.remove(oldKey) jsonObj.put(newKey, value) } } }) jsonArr.toJSONString }) // 调用UDF完成转换 val resultDf = df1.withColumn("nestedKey", when(col("nestedKey").isNotNull, replaceNestedJsonUdf(col("nestedKey"), typedLit(kMap), typedLit(vmap)) ).otherwise(col("nestedKey"))) resultDf.show(truncate = false)
两种方案输出结果均与你期望的结果完全一致。
内容的提问来源于stack exchange,提问作者Kundalini
相关产品推荐
相关产品推荐

