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

PySpark读取Kafka JSON序列化消息所有字段为NULL问题求助

问题排查与解决

核心原因分析

你的PySpark Schema定义与Kafka消息的JSON类型不匹配:

  • 原始JSON Schema中id、past_due、forborne被定义为number类型,JSON的number包含整数和浮点数两种类型
  • 你在PySpark中使用IntegerType接收这些字段,若Kafka消息里的数值是浮点数(比如123.0),from_json会因类型不匹配直接返回NULL

另外也需确认Kafka消息的JSON格式是否严格匹配Schema,比如是否存在字段名拼写错误、额外嵌套结构等问题。

解决方案

方案1:调整PySpark Schema类型兼容浮点数

将IntegerType替换为DoubleType,覆盖JSON number的类型范围:

from pyspark.sql.types import StructType, StructField, StringType, DoubleType

schema = StructType([
    StructField("id", DoubleType()),
    StructField("name", StringType()),
    StructField("past_due", DoubleType()),
    StructField("forborne", DoubleType())
])

如果确认消息中的数值都是整数,也可在解析后强制转换类型:

exploded_df = value_df.select(
    col("value.id").cast(IntegerType()).alias("id"),
    col("value.name").alias("name"),
    col("value.past_due").cast(IntegerType()).alias("past_due"),
    col("value.forborne").cast(IntegerType()).alias("forborne")
)

方案2:验证Kafka消息原始格式

临时修改代码打印原始消息,确认JSON结构和类型是否符合预期:

raw_query = streaming_df.select(col("value").cast("string").alias("raw_json")) \
    .writeStream \
    .format("console") \
    .outputMode("append") \
    .start()

raw_query.awaitTermination()

检查输出的raw_json是否为标准结构,比如:

{"id": 1, "name": "test_exposure", "past_due": 3, "forborne": 0}

若发现字段名大小写错误、额外嵌套等问题,需对应调整PySpark Schema。

方案3:自动推导Schema对比差异

临时用静态DataFrame让PySpark自动推导Schema,对比手动定义的Schema差异:

# 读取少量Kafka数据到静态DataFrame
static_df = spark.read \
    .format("kafka") \
    .options(**kafka_options) \
    .option("startingOffsets", "latest") \
    .load() \
    .select(col("value").cast("string").alias("json_str"))

# 自动推导Schema
inferred_schema = spark.read.json(static_df.rdd.map(lambda x: x.json_str)).schema
print(inferred_schema)

将推导结果与手动定义的Schema对比,可直接定位类型或结构不匹配的问题。

验证修改

调整Schema后重新运行流处理代码,即可正常解析字段值。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 00:32:36