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

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| ...
... +-----------------+-------------------+ ...

问题原因及修复方案

核心问题

  1. 列名大小写不匹配:Spark列名大小写敏感,代码中使用了c_Temperature(首字母大写),但实际列名是c_temperature(全小写),导致无法定位目标列返回NULL。
  2. 语法错误:withColumn("new_temperature", ...)一行缺少右括号,会导致代码执行报错。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 20:01:06