Spark DataFrame逐行动态替换JSON列字段的实现需求问询
解决Spark DataFrame动态JSON字段更新问题
嗨,我来帮你搞定这个逐行用变更JSON更新原始JSON的需求!这里有个实用的方案,基于Scala和json4s库来实现灵活的JSON合并,还能轻松支持数组类型的字段哦。
步骤1:需求核心拆解
我们要做的就是把每行chg字段里的键值对,覆盖到before字段的对应键上,同时保留before中未被修改的所有键值,最终生成after字段。
步骤2:实现JSON合并逻辑
因为JSON的键是动态可变的,用Spark内置的固定Schema函数不太方便,所以我们写一个自定义UDF(用户定义函数),结合json4s库来处理动态结构的JSON解析与合并。
首先,确保你的项目引入json4s-jackson依赖(如果是SBT项目,在build.sbt里添加):
libraryDependencies += "org.json4s" %% "json4s-jackson" % "4.0.6"
然后编写合并函数和注册UDF:
import spark.implicits._ import org.apache.spark.sql.functions._ import org.json4s._ import org.json4s.jackson.JsonMethods._ // 定义JSON合并逻辑:用chg的键值覆盖before对应键,保留其他键 def mergeJson(chgStr: String, beforeStr: String): String = { implicit val formats = DefaultFormats // 隐式转换规则,用于JSON与Scala对象互转 // 把JSON字符串解析成Scala Map val chgMap = parse(chgStr).extract[Map[String, Any]] val beforeMap = parse(beforeStr).extract[Map[String, Any]] // 合并Map:chgMap的键会覆盖beforeMap的同名键 val mergedMap = beforeMap ++ chgMap // 把合并后的Map重新序列化为JSON字符串 compact(render(mergedMap)) } // 注册为Spark可用的UDF val mergeJsonUdf = udf(mergeJson _)
步骤3:应用UDF生成目标DataFrame
把这个UDF应用到你的原始DataFrame上,就能得到想要的结果:
val df = Seq( (1, """{"b": "new", "c": "new"}""", """{"k": 1, "b": "old", "c": "old", "d": "old"}"""), (2, """{"b": "new", "d": "new"}""", """{"k": 2, "b": "old", "c": "old", "d": "old"}""") ).toDF("id", "chg", "before") // 生成结果DataFrame并筛选需要的字段 val resultDf = df.withColumn("after", mergeJsonUdf($"chg", $"before")) .select("id", "after") // 查看最终结果 resultDf.show(false)
执行后输出的结果完全符合你的预期:
+---+--------------------------------------------+ |id |after | +---+--------------------------------------------+ |1 |{"k":1,"b":"new","c":"new","d":"old"} | |2 |{"k":2,"b":"new","c":"old","d":"new"} | +---+--------------------------------------------+
关于数组类型的支持
这个方案天然支持数组字段!比如如果chg是{"arr": [1,2,3]},before里有{"arr": [4,5,6]},合并后after里的arr就会被替换成[1,2,3],完全满足你的需求。
内容的提问来源于stack exchange,提问作者Ged
相关产品推荐
相关产品推荐

