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本身返回字符串类型,不会影响内容,但可以删掉。
修复方案
- 先校验写入用的DataFrame正确性
写入前执行以下检查,避免用错数据集:
建议读取CSV时加参数去掉字段前后空格,规避表头尾随空格导致的隐藏问题:# 打印schema,确认只有value列,类型为string kafka_df.printSchema() # 取第一条数据查看原始字符串内容,规避show方法的转义干扰 print(kafka_df.take(1)[0]["value"])df = spark.read.format("csv") \ .option("header", "true") \ .option("ignoreLeadingWhiteSpace", "true") \ .option("ignoreTrailingWhiteSpace", "true") \ .load(f'{file_location}') - 根据使用的写入连接器修正代码
- 如果用开源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()
- 如果用开源Kafka连接器(包括连接Event Hubs的Kafka兼容端点),直接用转换好的
- 结果验证
写入完成后,消费时可以同时打印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
相关产品推荐
相关产品推荐

