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元素本身无关。另外你的后续处理逻辑绕了弯路,不需要把数组元素拼接成字符串,直接展开数组后提取结构体字段即可。
修正步骤
- 修正Schema:让StructField的名称和JSON中的key完全一致
- 简化处理流程:解析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
相关产品推荐
相关产品推荐

