Json creation from spark dataframe in scala

Viewed 23

Currently, we are converting a spark dataframe to JSON String to be sent to kafka.

In the process, we are doing toJSON twice which inserts \ for the inner json.

Snippet of the code:

val df=spark.sql("select * from dB.tbl")

val bus_dt="2022-09-23" 
case class kafkaMsg(busDate:String,msg:String)

Assuming my df has 2 columns as ID,STATUS, this will constitute the inner json of my kafka message.

JSON is created for msg and applied to case class.

val rdd=df.toJSON.rdd.map(msg=>kafkaMsg(busDate,msg))

Output at this step:

kafkaMsg(2022-09-23,{"id":1,"status":"active"})

Now, to send busDate and msg as JSON to kafka ,again a toJSON is applied.

val df1=spark.createDataFrame(rdd).toJSON

The output is:

{"busDate":"2022-09-23","msg":"{\"id\":1,\"status\":\"active\"}"}

The inner JSON is having \ which is not what the consumers are expecting.

Expected JSON:

{"busDate":"2022-09-23","msg":{"id":1,"status":"active"}}

How can I create this json without \ and send to kafka.

Please note the msg value varies and cannot be mapped to a schema.

0 Answers
Related