Scala中Spark DataFrame转JSON写Kafka的转义符问题求助
解决Spark DataFrame转JSON发送Kafka时内部JSON转义问题
问题背景
将Spark DataFrame转换为JSON发送至Kafka时,因两次调用toJSON导致内部JSON出现转义符\。原流程先将每行数据转成JSON字符串存入case class的msg字段,再将case class转成JSON,最终输出的msg字段是带转义的字符串而非嵌套JSON对象,不符合消费者预期。且msg字段结构不固定,无法预先映射到固定schema。
解决方案
方法1:利用Spark结构化API构造嵌套结构(推荐)
直接将原始DataFrame的所有列打包成struct类型,再与busDate组合后转JSON,避免将msg转为字符串再嵌套的操作。
import org.apache.spark.sql.functions._ val df = spark.sql("select * from dB.tbl") val bus_dt = "2022-09-23" // 动态将原始DF的所有列打包为struct,命名为msg val nestedDF = df .withColumn("busDate", lit(bus_dt)) // 添加busDate字段 .select(col("busDate"), struct(df.columns.map(col): _*).alias("msg")) // 构造嵌套结构 // 直接转JSON,得到无转义的嵌套结果 val kafkaJsonDF = nestedDF.toJSON
此方法优势:
- 借助Spark内置函数自动处理JSON结构,无需手动拼接,避免语法错误
- 动态适配
msg字段的任意结构,无需修改代码适配schema变化 - 输出结果符合预期:
{"busDate":"2022-09-23","msg":{"id":1,"status":"active"}}
方法2:手动构造JSON字符串
如果需要更灵活的JSON结构控制,可以直接在RDD层面手动拼接最终的JSON字符串,跳过两次toJSON的步骤。
val df = spark.sql("select * from dB.tbl") val bus_dt = "2022-09-23" // 手动拼接嵌套JSON val kafkaJsonRDD = df.rdd.map { row => val msgJson = row.toJSON // 获取每行的原始JSON(无外层引号) s"""{"busDate":"$bus_dt","msg":$msgJson}""" } // 转为DataFrame用于后续Kafka发送 val kafkaJsonDF = spark.createDataFrame(kafkaJsonRDD.map(Tuple1.apply)).toDF("value")
注意:如果bus_dt包含特殊字符(如双引号),需要额外做转义处理,避免JSON格式错误。
内容的提问来源于stack exchange,提问作者AMRITA NEHA
相关产品推荐
相关产品推荐

