Scala中DataFrame满足FLAG=Y时合并两列JSON的实现方法
错误原因
scala.util.parsing.json.JSON.parseFull是Scala标准库中已标记废弃的JSON解析工具,容错性极差,只要遇到空字符串、格式不标准的JSON(比如尾逗号、未转义特殊字符、多余空格)就会返回None,代码中直接调用.get取值时,只要任意一行解析失败就会抛出None.get对应的空值错误。同时该逐行硬解析的方式无法利用Spark内置的JSON处理优化,性能远低于原生函数方案。
单语句实现方案
直接使用Spark内置JSON处理函数实现,无需手动解析Map,天然满足「JSON2重复键优先级更高」的合并规则(map_concat合并时,后传入Map的重复键会覆盖前序Map的同key值),核心实现代码如下:
import org.apache.spark.sql.functions._ val resDf = sourceDf.withColumn( "JSON2", when( col("FLAG") === "Y", to_json( map_concat( from_json(col("JSON1"), lit("map<string,string>")), from_json(col("JSON2"), lit("map<string,string>")) ) ) ).otherwise(col("JSON2")) )
逻辑验证:输入样例为FLAG=Y, JSON1={"v1":"1","v2":"2","v3":"5"}, JSON2={"v1":"3","v2":"4","v4":"6"}时,执行后JSON2的输出为{"v1":"3","v2":"4","v3":"5","v4":"6"},和预期结果完全匹配。
注意事项
- 若JSON值存在嵌套结构、数字/布尔类型等非字符串值,可将
from_json的schema从map<string,string>调整为map<string,json>,合并逻辑无需改动 - 若数据中存在脏JSON记录,可给
from_json增加容错配置,避免单条坏数据导致任务中断,示例写法:from_json( col("JSON1"), lit("map<string,string>"), Map("mode" -> "PERMISSIVE") ) - 若JSON结构为固定字段的结构体而非动态键Map,可将schema替换为对应StructType,
map_concat替换为struct函数按字段选取即可,优先级规则保持JSON2字段覆盖JSON1同名字段。
内容的提问来源于stack exchange,提问作者SD'Anc
相关产品推荐
相关产品推荐

