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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 05:06:18