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

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对象会更简单:

  1. 先引入JSON处理依赖(以FastJSON为例,集群环境已导入可跳过)
<dependency>
    <groupId>com.alibaba</groupId>
    <artifactId>fastjson</artifactId>
    <version>1.2.83</version>
</dependency>
  1. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 00:06:00