PySpark DataFrame展开JSON数组列失败,请求排查错误原因
问题
我在PySpark DataFrame中有一个名为substitutions的JSON字符串列,该列存储着包含多个对象的数组,希望将数组展开为每行对应一个元素的形式,DataFrame中还有其他列。
DataFrame结构如下:
+--------------------+----------+--------------------+--------------------+----------+--------------------+---------+--------+--------------------+ | requestid|sourcepage| cartid| tm| dt| customerId| usItemId|prefType| substitutions| +--------------------+----------+--------------------+--------------------+----------+--------------------+---------+--------+--------------------+ |00-efbedfe05b4482...| CHECKOUT|808b44cc-1a38-4dd...|2023-04-25 00:07:...|2023-04-25|f1a34e16-a6d0-6f5...|862776084| NO_PREF|{"id":{"productId...| +--------------------+----------+--------------------+--------------------+----------+--------------------+---------+--------+--------------------+
substitutions列的JSON字符串内容为:
[ { "id": { "productId": "2N3UYGUTROQK", "usItemId": "32667929" }, "usItemId": "32667929", "itemRank": 1, "customerChoice": false }, { "id": { "productId": "2N3UYGUTRHQK", "usItemId": "32667429" }, "usItemId": "32667429", "itemRank": 2, "customerChoice": true }, { "id": { "productId": "2N3UYGUTRYQK", "usItemId": "32667529" }, "usItemId": "32667529", "itemRank": 3, "customerChoice": false }, { "id": { "productId": "2N3UYGUTIQK", "usItemId": "32667329" }, "usItemId": "32667329", "itemRank": 4, "customerChoice": false }, { "id": { "productId": "2N3UYGUTYOQK", "usItemId": "32663929" }, "usItemId": "32663929", "itemRank": 5, "customerChoice": false } ]
尝试了以下代码但未得到预期结果:
df.select("*", f.explode(f.from_json("substitutions", MapType(StringType(),StringType()))))
结果:
+--------------------+----------+--------------------+--------------------+----------+--------------------+---------+--------+--------------------+-------+ | requestid|sourcepage| cartid| tm| dt| customerId| usItemId|prefType| substitutions|entries| +--------------------+----------+--------------------+--------------------+----------+--------------------+---------+--------+--------------------+-------+ |00-efbedfe05b4482...| CHECKOUT|808b44cc-1a38-4dd...|2023-04-25 00:07:...|2023-04-25|f1a34e16-a6d0-6f5...|862776084| NO_PREF|[{"id":{"productI...| null| +--------------------+----------+--------------------+--------------------+----------+--------------------+---------+--------+--------------------+-------+
请问我哪里出错了?
错误原因
你用错了Schema类型:
substitutions是JSON数组,里面每个元素是嵌套结构的对象,不是简单的键值对Map。MapType(StringType(), StringType())只能解析键值都是字符串的JSON对象,无法匹配数组+嵌套对象的结构,导致from_json返回null,后续explode自然也得不到有效数据。
正确解决方案
需要先定义匹配JSON结构的Schema,解析数组后再展开:
步骤1:定义Schema
根据substitutions的JSON结构,用ArrayType包裹StructType,包含所有嵌套字段:
from pyspark.sql.types import ( StructType, StructField, StringType, IntegerType, BooleanType, ArrayType ) # 定义嵌套的id结构 id_schema = StructType([ StructField("productId", StringType()), StructField("usItemId", StringType()) ]) # 定义数组元素的结构 substitution_schema = StructType([ StructField("id", id_schema), StructField("usItemId", StringType()), StructField("itemRank", IntegerType()), StructField("customerChoice", BooleanType()) ]) # 最终的数组Schema full_schema = ArrayType(substitution_schema)
步骤2:解析JSON并展开数组
先解析substitutions列成数组类型,再用explode展开为多行:
from pyspark.sql import functions as f # 解析JSON列 df_parsed = df.withColumn("substitutions_parsed", f.from_json("substitutions", full_schema)) # 展开数组,同时保留原有的其他列 df_exploded = df_parsed.select("*", f.explode("substitutions_parsed").alias("substitution")) # 可选:展开嵌套字段,把substitution里的字段拆成单独列 df_final = df_exploded.drop("substitutions", "substitutions_parsed") \ .select( "requestid", "sourcepage", "cartid", "tm", "dt", "customerId", "usItemId", "prefType", "substitution.id.productId", "substitution.id.usItemId", "substitution.usItemId", "substitution.itemRank", "substitution.customerChoice" )
最终效果
展开后每行对应一个substitution元素,原有的其他列会被复制到每行,嵌套的字段也可以单独提取出来。
内容的提问来源于stack exchange,提问作者Shibu
相关产品推荐
相关产品推荐

