PySpark读取Kafka流时Value字段JSON转换及转义符消除问题
解决Glue Spark读取Kafka流后JSON输出含转义符的问题
问题场景
当前使用Glue通过Spark读取Kafka流,核心代码如下:
try: options = { "kafka.sasl.jaas.config": 'org.apache.kafka.common.security.plain.PlainLoginModule required username="XXXXXXXXXXXX" password="XXXXXXXXXXXXXX";', "kafka.sasl.mechanism": "PLAIN", "kafka.security.protocol": "SASL_SSL", "kafka.bootstrap.servers": "kafka-server:9092", "subscribe": "masterstaging_cfr_out_customeragreement_event_disbursement_ini", "startingOffsets":"latest" } df = spark.readStream.format("kafka").options(**options).load() df=df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") df.writeStream.format("json") \ .option("checkpointLocation", "s3://output/")\ .outputMode("append") \ .option("path", "s3://output/") \ .start() \ .awaitTermination() except Exception as e: print(e)
输出的JSON中value字段带有大量转义符\,示例如下:
{ "key": "test", "value": "{\n \"payload\": {\n \"EventCode\": {\n \"operation_code\": \"Creation\",\n \"reason_code\": \"\"\n },\n \"Data\": \n \"id\": 8888881,\n \"ref\": \"D16/0405\" \n }\n }\n}" }
需要修改df=df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")这部分代码,消除输出中的转义符。
解决方案
问题根源是仅将value转换为字符串,Spark写入JSON时会把该字符串当作普通文本处理,自动添加转义符。正确做法是将value解析为JSON结构体,让Spark识别其嵌套结构。
方法1:定义Schema解析(推荐)
- 先根据
value的JSON结构定义对应的Spark Schema:
from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 匹配你的value字段JSON结构 schema = StructType([ StructField("payload", StructType([ StructField("EventCode", StructType([ StructField("operation_code", StringType()), StructField("reason_code", StringType()) ])), StructField("Data", StructType([ StructField("id", IntegerType()), StructField("ref", StringType()) ])) ])) ])
- 替换原有的selectExpr代码,用
from_json解析字符串为JSON结构体:
# 替换原来的df=df.selectExpr(...) df = df.select( df.key.cast("string").alias("key"), from_json(df.value.cast("string"), schema).alias("value") )
方法2:使用字符串形式Schema(适合快速测试)
如果不想单独定义Schema对象,也可以直接在selectExpr中传入JSON格式的Schema字符串:
df = df.selectExpr( "CAST(key AS STRING)", "from_json(CAST(value AS STRING), '{\"payload\": {\"EventCode\": {\"operation_code\": \"string\", \"reason_code\": \"string\"}, \"Data\": {\"id\": \"integer\", \"ref\": \"string\"}}}') AS value" )
效果验证
修改后输出的JSON会是嵌套结构,无转义符,示例如下:
{ "key": "test", "value": { "payload": { "EventCode": { "operation_code": "Creation", "reason_code": "" }, "Data": { "id": 8888881, "ref": "D16/0405" } } } }
可选:扁平化输出
如果不需要保留外层的value字段,可直接展开嵌套字段:
df = df.select( df.key.cast("string").alias("key"), from_json(df.value.cast("string"), schema).payload.EventCode.operation_code.alias("operation_code"), from_json(df.value.cast("string"), schema).payload.EventCode.reason_code.alias("reason_code"), from_json(df.value.cast("string"), schema).payload.Data.id.alias("id"), from_json(df.value.cast("string"), schema).payload.Data.ref.alias("ref") )
输出会变成扁平化的JSON:
{ "key": "test", "operation_code": "Creation", "reason_code": "", "id": 8888881, "ref": "D16/0405" }
内容的提问来源于stack exchange,提问作者Smaillns
相关产品推荐
相关产品推荐

