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。
排查及解决步骤
- 先确认原始数据格式,输出前10条原始Body校验结构:
display(mess.limit(10))
- 如果确认是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.*")
- 如果是字符串格式数字的问题,可以开启
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
相关产品推荐
相关产品推荐

