PySpark中get_json_object仅返回NULL值问题排查求助
解决Spark提取JSON字符串字段返回NULL的问题
问题场景
在Palantir Foundry中,有一个Spark DataFrame,其中c_temperature列存储的是如下格式的JSON字符串:
{"TempCelsiusEndAvg":"24.33","TempCelsiusEndMax":"null","TempCelsiusEndMin":"null","TempCelsiusStartAvg":"22.54","TempCelsiusStartMax":"null","TempCelsiusStartMin":"null","TempEndPlausibility":"T_PLAUSIBLE","TempStartPlausibility":"T_PLAUSIBLE"}
想要提取其中的TempCelsiusEndAvg和TempCelsiusStartAvg字段生成新列,尝试了以下代码,但新列全为NULL:
from pyspark.sql import functions as F from pyspark.sql.types import StringType def flat_json(sessions_finished): df = sessions_finished df = df.withColumn("new_temperature", F.col('c_temperature').cast(StringType()) df = df.withColumn("TempCelsiusEndAvg", F.get_json_object("c_Temperature", '$.TempCelsiusEndAvg')) df = df.withColumn("TempCelsiusStartAvg", F.get_json_object("c_Temperature", '$.TempCelsiusStartAvg')) return df
期望结果
... +-----------------+-------------------+ ... ... |TempCelsiusEndAvg|TempCelsiusStartAvg| ... ... +-----------------+-------------------+ ... ... | 24.33| 22.54| ... ... +-----------------+-------------------+ ... ... | 29.28| 25.16| ... ... +-----------------+-------------------+ ... ... | null| null| ... ... +-----------------+-------------------+ ...
实际结果
... +-----------------+-------------------+ ... ... |TempCelsiusEndAvg|TempCelsiusStartAvg| ... ... +-----------------+-------------------+ ... ... | null| null| ... ... +-----------------+-------------------+ ... ... | null| null| ... ... +-----------------+-------------------+ ... ... | null| null| ... ... +-----------------+-------------------+ ...
问题原因及修复方案
核心问题
- 列名大小写不匹配:Spark列名大小写敏感,代码中使用了
c_Temperature(首字母大写),但实际列名是c_temperature(全小写),导致无法定位目标列返回NULL。 - 语法错误:
withColumn("new_temperature", ...)一行缺少右括号,会导致代码执行报错。 - JSON字符串"null"未处理:原JSON中的"null"是字符串类型,直接提取后会保留字符串值,需要转为Spark原生的NULL值。
修正后的代码
from pyspark.sql import functions as F from pyspark.sql.types import StringType, DoubleType def flat_json(sessions_finished): df = sessions_finished # 确保列是字符串类型(若原列类型非字符串) df = df.withColumn("c_temperature_str", F.col('c_temperature').cast(StringType())) # 提取JSON字段,注意列名大小写正确 df = df.withColumn("TempCelsiusEndAvg", F.get_json_object("c_temperature_str", '$.TempCelsiusEndAvg')) df = df.withColumn("TempCelsiusStartAvg", F.get_json_object("c_temperature_str", '$.TempCelsiusStartAvg')) # 将字符串"null"转为Spark NULL,并转为浮点类型 df = df.withColumn("TempCelsiusEndAvg", F.when(F.col("TempCelsiusEndAvg") != "null", F.col("TempCelsiusEndAvg").cast(DoubleType())) .otherwise(F.lit(None))) df = df.withColumn("TempCelsiusStartAvg", F.when(F.col("TempCelsiusStartAvg") != "null", F.col("TempCelsiusStartAvg").cast(DoubleType())) .otherwise(F.lit(None))) # 可选:删除中间辅助列 df = df.drop("c_temperature_str") return df
更高效的批量解析方案(使用from_json)
如果需要提取多个JSON字段,结合Schema解析的方式更易维护:
from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, StringType, DoubleType def flat_json(sessions_finished): # 定义JSON结构Schema temp_schema = StructType([ StructField("TempCelsiusEndAvg", StringType()), StructField("TempCelsiusStartAvg", StringType()), StructField("TempEndPlausibility", StringType()), StructField("TempStartPlausibility", StringType()) ]) df = sessions_finished # 解析JSON字符串为结构体 df = df.withColumn("temp_struct", F.from_json(F.col('c_temperature').cast(StringType()), temp_schema)) # 提取字段并处理字符串"null" df = df.withColumn("TempCelsiusEndAvg", F.when(F.col("temp_struct.TempCelsiusEndAvg") != "null", F.col("temp_struct.TempCelsiusEndAvg").cast(DoubleType())) .otherwise(F.lit(None))) df = df.withColumn("TempCelsiusStartAvg", F.when(F.col("temp_struct.TempCelsiusStartAvg") != "null", F.col("temp_struct.TempCelsiusStartAvg").cast(DoubleType())) .otherwise(F.lit(None))) # 可选:删除中间结构体列 df = df.drop("temp_struct") return df
内容的提问来源于stack exchange,提问作者steffsen
相关产品推荐
相关产品推荐

