PySpark中安全访问不存在的嵌套JSON属性问题
问题描述
读取JSON文件后,我尝试用以下PySpark代码创建列:
observation_df.withColumn("contained_observations", F.explode(col("contained"))) .withColumn("code", col("contained_observations.code")) .withColumn("code_text", col("code.text")) .withColumn("coding", when(col("code").isNotNull(), col("code").getField("coding")).otherwise(None)) .select( col("code"), col("code_text") # col("coding") ) .printSchema()
但code结构体中不存在coding字段,即使加了判断,仍触发错误:
AnalysisException: No such struct field coding in text
我希望即使输入中没有该字段,也能将其加入DataFrame并赋值为None,该如何实现?
我还尝试了以下两种写法:
.withColumn("coding", when(col("code").isNotNull(), col("code").getField("coding")).otherwise(None)) .withColumn("coding", col("code").getField("coding").isNotNull())
读取JSON时未指定Schema(Schema不固定且预先未知),由Spark自动推断,当前Schema为:
root -> code(struct) -> text(String)
解决方案
问题核心是Spark在逻辑计划阶段就会校验结构体字段是否存在,不会等到运行时判断,所以直接用getField会触发解析错误。以下几种方法可以实现需求:
方法1:使用try_get_field(Spark 3.3+ 推荐)
Spark 3.3及以上版本提供了try_get_field函数,专门处理不确定字段是否存在的场景——字段存在则返回对应值,不存在则返回null:
from pyspark.sql import functions as F observation_df.withColumn("contained_observations", F.explode(col("contained"))) .withColumn("code", col("contained_observations.code")) .withColumn("code_text", col("code.text")) .withColumn("coding", F.try_get_field(col("code"), "coding")) .select("code", "code_text", "coding") .printSchema()
方法2:自定义UDF兼容低版本Spark
如果你的Spark版本低于3.3,可自定义UDF检查结构体字段:
from pyspark.sql import functions as F from pyspark.sql.types import StringType def get_coding_field(code_struct): if code_struct and isinstance(code_struct, dict) and "coding" in code_struct: return code_struct["coding"] return None # 可根据实际场景修改返回类型(如ArrayType、StructType等) get_coding_udf = F.udf(get_coding_field, StringType()) observation_df.withColumn("contained_observations", F.explode(col("contained"))) .withColumn("code", col("contained_observations.code")) .withColumn("code_text", col("code.text")) .withColumn("coding", get_coding_udf(col("code"))) .select("code", "code_text", "coding") .printSchema()
方法3:动态扩展Schema(通用方案)
先获取当前code字段的Schema,动态添加coding字段后重新解析,后续即可安全引用:
from pyspark.sql import functions as F from pyspark.sql.types import StructField, NullType # 获取当前code字段的Schema结构 current_code_schema = observation_df.select("contained_observations.code").schema["code"].dataType # 检查并添加coding字段(若不存在) if "coding" not in [field.name for field in current_code_schema.fields]: new_code_schema = StructType( current_code_schema.fields + [StructField("coding", NullType(), nullable=True)] ) observation_df = observation_df.withColumn( "code", F.from_json(F.to_json(col("code")), new_code_schema) ) # 后续可直接引用code.coding observation_df.withColumn("contained_observations", F.explode(col("contained"))) .withColumn("code", col("contained_observations.code")) .withColumn("code_text", col("code.text")) .withColumn("coding", col("code.coding")) .select("code", "code_text", "coding") .printSchema()
内容的提问来源于stack exchange,提问作者Shreyas Patil
相关产品推荐
相关产品推荐

