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

Spark写入Kafka的JSON数据读取时仅返回id字段异常排查

问题原因

异常根源在写入环节,实际存入Kafka/Event Hub的消息体只有id字段值,和你预期的完整JSON不符,和你看到的转义引号无关——show方法打印字符串时会自动给字符串加外层双引号,内部双引号用两个双引号转义,属于正常显示效果,不是数据错误。
常见触发原因有两个:

  • 写入时用错DataFrame:你预览打印的是经过JSON转换的kafka_df,但实际调用write接口时传入的是读取CSV得到的原始df。如果你用的是Azure Event Hubs原生Spark连接器(而非开源Kafka兼容API),该连接器要求消息体必须是名为body的二进制列,部分旧版本找不到该列时不会直接抛错,会默认取DataFrame第一列(也就是id列)作为消息体发送,最终存到主题里的内容就只有UUID格式的id值。
  • 写入配置错误:写入时额外配置了消息值映射规则,比如强制指定id列作为value/body发送,覆盖了你提前构建好的包含完整JSON的value列。
    另外你代码里的.selectExpr("CAST(value AS STRING)")是冗余操作,to_json本身返回字符串类型,不会影响内容,但可以删掉。
修复方案
  1. 先校验写入用的DataFrame正确性
    写入前执行以下检查,避免用错数据集:
    # 打印schema,确认只有value列,类型为string
    kafka_df.printSchema()
    # 取第一条数据查看原始字符串内容,规避show方法的转义干扰
    print(kafka_df.take(1)[0]["value"])
    
    建议读取CSV时加参数去掉字段前后空格,规避表头尾随空格导致的隐藏问题:
    df = spark.read.format("csv") \
        .option("header", "true") \
        .option("ignoreLeadingWhiteSpace", "true") \
        .option("ignoreTrailingWhiteSpace", "true") \
        .load(f'{file_location}')
    
  2. 根据使用的写入连接器修正代码
    • 如果用开源Kafka连接器(包括连接Event Hubs的Kafka兼容端点),直接用转换好的kafka_df写入即可,不要额外指定value列映射:
      kafka_df.write \
          .format("kafka") \
          .option("kafka.bootstrap.servers", "你的服务地址") \
          .option("topic", "你的主题名") \
          # 若使用Event Hubs Kafka API需要补充以下SASL配置,原生Kafka不需要
          .option("kafka.sasl.mechanism", "PLAIN") \
          .option("kafka.security.protocol", "SASL_SSL") \
          .option("kafka.sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username='$ConnectionString' password='你的Event Hubs连接字符串';") \
          .save()
      
    • 如果用Azure Event Hubs原生连接器,需要将value列重命名为body并转为二进制类型再写入:
      eh_write_df = kafka_df.select(F.col("value").cast("binary").alias("body"))
      eh_write_df.write \
          .format("eventhubs") \
          .option("connectionString", "你的Event Hubs连接字符串;EntityPath=你的主题名") \
          .save()
      
  3. 结果验证
    写入完成后,消费时可以同时打印key、value字段确认内容:
    string_df = df.select(
        F.col("key").cast("string").alias("key"),
        F.col("value").cast("string").alias("value")
    )
    string_df.show(truncate=False)
    

内容的提问来源于stack exchange,提问作者user14681827

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 00:27:22