PySpark如何判断单条记录是否包含指定JSON字段
问题说明
- 运行环境:PySpark 3.1.1、Python 3.8
现有输入JSON文件内容:
{"id" :1, "field_a": "test"} {"id" :2, "field_a": "test", "field_b": "z"} {"id" :3, "field_a": "test", "field_b": null}
使用Spark默认JSON接口读取文件时,会自动为缺失对应字段的记录填充null值,读取代码如下:
spark = SparkSession \ .builder \ .appName("Python Spark SQL basic example") \ .getOrCreate() df = spark.read.json('my_file.json') df.show()
默认读取返回结果:
+-------+-------+---+ |field_a|field_b| id| +-------+-------+---+ | test| null| 1| | test| z| 2| | test| null| 3| +-------+-------+---+
该默认机制存在缺陷:无法区分字段值为NULL和记录本身不存在对应字段两种场景。内置方法df.columns仅能获取全局表结构的列信息,无法判断单条记录内的字段存在性。
预期实现record_has_column(field_b)函数,通过df = df.withColumn("column_in_json", record_has_column(field_b))调用后,输出如下结果:字段本身不存在时返回false,字段存在(即使值为null)时返回true。
+-------+-------+---+--------------+ |field_a|field_b| id|column_in_json| +-------+-------+---+--------------+ | test| null| 1| false| | test| z| 2| true| | test| null| 3| true| +-------+-------+---+--------------+
实现方法
Spark默认JSON解析器会在Schema合并阶段直接抹除「字段缺失」和「字段值为null」的差异,因此需要在读取阶段保留原始JSON结构,通过Map类型解析判断字段是否存在:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.types import MapType, StringType spark = SparkSession \ .builder \ .appName("Python Spark SQL basic example") \ .getOrCreate() # 第一步:以文本格式读取原始JSON行,不做自动结构化解析 df_raw = spark.read.text("my_file.json") # 第二步:将每行JSON字符串解析为Map结构,key为字段名,value为字段值 df = df_raw.withColumn( "json_data", F.from_json(F.col("value"), MapType(StringType(), StringType())) ) # 第三步:展开需要的字段,同时通过Map的key判断字段是否存在 result = df.select( F.col("json_data")["field_a"].alias("field_a"), F.col("json_data")["field_b"].alias("field_b"), F.col("json_data")["id"].cast("int").alias("id"), # 核心判断逻辑:Map的key列表是否包含目标字段名 F.map_keys("json_data").contains("field_b").alias("column_in_json") ) result.show()
如果需要封装为通用的record_has_column函数,可以参考如下实现:
def record_has_column(col_name: str): return F.map_keys("json_data").contains(col_name) # 调用示例 result = result.withColumn("check_field_a", record_has_column("field_a"))
运行后输出结果和预期完全一致:id=1的记录因原始JSON无field_b字段返回false,id=3的记录虽field_b值为null但字段存在,返回true。
内容的提问来源于stack exchange,提问作者Hyruma92
相关产品推荐
相关产品推荐

