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

Pyspark调用from_json转换EventHub流式数据后出现null值问题

PySpark from_json解析EventHub JSON整数字段返回null问题原因及解决方案

问题核心原因

Spark from_json 函数默认使用严格模式解析数据,类型不匹配、数据结构不匹配、键名不匹配时都会直接返回null,不会抛出异常,这是你将字段改为StringType后解析正常的核心原因,具体可以排查以下几个场景:

  • JSON结构不匹配:你的Azure函数返回的是Vehicles数组序列化后的JSON字符串,即EventHub的Body是[{}, {}, ...]格式的JSON数组,而你当前定义的Schema是单个JSON对象的结构,直接解析数组自然会返回全量null。
  • 数值字段带引号:源API返回的部分整数字段实际是字符串格式(比如"Line": "12"而非"Line":12),Spark默认不会自动将字符串格式的数字强转为IntegerType,直接解析会返回null。
  • 字段存在非预期值:部分行的整数字段存在空字符串""、浮点数值(比如"Delay": 2.5)、非数字字符等不符合IntegerType定义的值,解析会返回null。
  • 键名不匹配:需要确认源JSON的键名和你Schema中定义的完全一致,包括大小写、空格(比如你定义的GPS Quality带空格,要确认源数据键名完全相同),键名不匹配时对应字段会返回null。

排查及解决步骤

  1. 先确认原始数据格式,输出前10条原始Body校验结构:
display(mess.limit(10))
  1. 如果确认是JSON数组结构,修改解析逻辑先解析为数组再展开:
from pyspark.sql.functions import from_json, col, explode
from pyspark.sql.types import ArrayType

# 将原有Schema包装为数组类型
array_schema = ArrayType(schema)
df = mess.select(from_json(col("Body"), array_schema).alias("data_arr")) \
         .select(explode("data_arr").alias("data")) \
         .select("data.*")
  1. 如果是字符串格式数字的问题,可以开启from_json宽松模式自动转换,或者先读为StringType再强转:

方式1:开启宽松解析模式

df = mess.select(
    from_json(
        col("Body"), 
        array_schema, # 若为单个对象结构直接用原schema即可
        {"mode": "PERMISSIVE"}
    ).alias("data_arr")
)

方式2:先读为StringType再强转,方便定位异常数据

# 先把整数字段都定义为StringType解析
string_schema = StructType([
StructField("DataGenerated", TimestampType()),
StructField("Line", StringType()),
StructField("Route", StringType()),
StructField("VehicleCode", StringType()),
StructField("VehicleService", StringType()),
StructField("VehicleID", StringType()),
StructField("Speed", StringType()),
StructField("Delay", StringType()),
StructField("Lat", DoubleType()),
StructField("Lon", DoubleType()),
StructField("GPS Quality", StringType()),
])

# 解析后再强转为int,用try_cast避免转换失败报错
from pyspark.sql.functions import try_cast
df = mess.select(from_json(col("Body"), ArrayType(string_schema)).alias("data_arr")) \
         .select(explode("data_arr").alias("data")) \
         .select("data.*") \
         .withColumn("Line", try_cast(col("Line"), IntegerType())) \
         .withColumn("VehicleCode", try_cast(col("VehicleCode"), IntegerType()))

# 可以筛选出转换失败的行排查异常值
display(df.filter(col("Line").isNull() & col("Line").isNotNull()))

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 07:18:01