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

PySpark中from_json解析含Null的JSON时全字段为空问题求助

PySpark处理含Null的JSON字符串转DataFrame问题

问题场景

读取JSON格式的字符串(注意addresses是字符串类型,不是数组)转成DataFrame多列时,所有元素非空时正常,但只要有Null元素,所有字段都会变成Null。

输入数据

id='1'
addresses='
[{
    "city": "city1",
    "state": null,
    "street": null,
    "postalCode": null,
    "country": "country1"
}
,
{
    "city": "city2",
    "state": null,
    "street": "street2",
    "postalCode": "11111",
    "country": "country2"
}]'

期望输出

id  city    state   street  postalCode  country
1   city1   null    null    null        country1
1   city2   null    street2 11111       country2

现有代码

addl_addr_schema = ArrayType(StructType([
    StructField("addl_addr_city", StringType(), True),
    StructField("addl_addr_state", StringType(), True),
    StructField("addl_addr_street", StringType(), True),
    StructField("addl_addr_postalCode", StringType(), True),
    StructField("addl_addr_country", StringType(), True),
]))


dpDF_transformed = dpDF_temp.withColumn('addresses_transformed', from_json('addresses', addl_addr_schema)) \
                                    .withColumn('addl_addr', explode_outer('addresses_transformed'))

dpDF_transformed = dpDF_transformed.select("*",col("addresses_transformed.addl_addr_street").alias("addl_addr_street_array"),col("addresses_transformed.addl_addr_city").alias("addl_addr_city_array"),col("addresses_transformed.addl_addr_state").alias("addl_addr_state_array"),col("addresses_transformed.addl_addr_postalCode").alias("addl_addr_postalCode_array"),col("addresses_transformed.addl_addr_country").alias("addl_addr_country_array"))

dpDF_final = dpDF_transformed.withColumn("addl_addr_street",concat_ws(",","addl_addr_street_array")) \
                                     .withColumn("addl_addr_city",concat_ws(",","addl_addr_city_array")) \
                .withColumn("addl_addr_state",concat_ws(",","addl_addr_state_array")) \
.withColumn("addl_addr_postalCode",concat_ws(",","addl_addr_postalCode_array")) \
                                     .withColumn("addl_addr_country",concat_ws(",","addl_addr_country_array")) \
                                    .drop("addresses","addresses_transformed","addl_addr","addl_addr_street_array","addl_addr_city_array","addl_addr_state_array","addl_addr_postalCode_array","addl_addr_country_array")

实际输出

id  city    state   street  postalCode  country
1   city1   null    null    null        null
1   city2   null    null    null        null

问题原因及解决方案

核心问题

你定义的Schema字段名和JSON里的字段名完全不匹配!JSON里的字段是city、state,但你Schema里写的是addl_addr_city、addl_addr_state,这才导致from_json解析失败,最终所有字段变成Null,和Null元素本身无关。另外你的后续处理逻辑绕了弯路,不需要把数组元素拼接成字符串,直接展开数组后提取结构体字段即可。

修正步骤

  1. 修正Schema:让StructField的名称和JSON中的key完全一致
  2. 简化处理流程:解析JSON字符串为数组后,直接展开数组,再提取结构体中的各个字段

修正后的完整代码

from pyspark.sql import SparkSession
from pyspark.sql.types import ArrayType, StructType, StructField, StringType
from pyspark.sql.functions import from_json, explode_outer, col

# 初始化SparkSession(如果还没初始化)
spark = SparkSession.builder.appName("JSONtoDF").getOrCreate()

# 构造测试DataFrame
dpDF_temp = spark.createDataFrame([(1, '''
[{
    "city": "city1",
    "state": null,
    "street": null,
    "postalCode": null,
    "country": "country1"
}
,
{
    "city": "city2",
    "state": null,
    "street": "street2",
    "postalCode": "11111",
    "country": "country2"
}]''')], ["id", "addresses"])

# 定义正确的Schema:字段名和JSON中的key完全匹配
addr_schema = ArrayType(StructType([
    StructField("city", StringType(), nullable=True),
    StructField("state", StringType(), nullable=True),
    StructField("street", StringType(), nullable=True),
    StructField("postalCode", StringType(), nullable=True),
    StructField("country", StringType(), nullable=True),
]))

# 1. 解析JSON字符串为数组列
# 2. 展开数组,每个地址生成一行
dpDF_transformed = dpDF_temp.withColumn('addresses_array', from_json('addresses', addr_schema)) \
                           .withColumn('addr', explode_outer('addresses_array'))

# 提取结构体中的各个字段,组成最终的DataFrame
dpDF_final = dpDF_transformed.select(
    "id",
    col("addr.city"),
    col("addr.state"),
    col("addr.street"),
    col("addr.postalCode"),
    col("addr.country")
).drop("addresses", "addresses_array", "addr")

# 查看结果
dpDF_final.show()

输出验证

运行上述代码后,输出会和你期望的完全一致,Null元素会被正确保留,不会导致其他字段变成Null。

内容的提问来源于stack exchange,提问作者Jatin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 10:30:52