You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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解析(推荐)

  1. 先根据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())
        ]))
    ]))
])
  1. 替换原有的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.18 02:45:36