PySpark DataFrame字符串列中JSON数组的展开与解析
处理PySpark DataFrame中JSON数组字符串的解析与行扩展
问题场景
我有一个PySpark DataFrame,其中mappingresult列是字符串格式,内部包含JSON数组(部分行的该列为空数组[])。
原始DataFrame创建代码
spark.createDataFrame(pd.DataFrame({'server': {0: '3456gj', 1: '56ujdn98', 2: '56v95bd', 3: 'barca6mw2k'}, 'logev': {0: 'map.map', 1: 'map.map', 2: 'map.map', 3: 'map.map'}, 'status': {0: 'no prod', 1: 'as rules', 2: 'is found', 3: 'not found'}, 'mappingresult': {0: '[]', 1: '[]', 2: '[{\"model\":\"s\",\"com\":\"48b4-bfde\",\"ctid\":987456,\"params\":[{\"ID\":\"tr\",\"val\":\"399.00\"},{\"IDp\":\"merch\",\"val\":\"stackoverflow\"}]},{\"model\":\"s\",\"com\":\"76r-bfde\",\"ctid\":987456,\"params\":[{\"ID\":\"tr\",\"val\":\"399.00\"},{\"IDp\":\"merch\",\"val\":\"stackoverflow\"}]}]', 3: '[{\"model\":\"s\",\"com\":\"4827-44e9\",\"ctid\":987456,\"params\":[{\"ID\":\"tr\",\"val\":\"399.00\"},{\"IDp\":\"merch\",\"val\":\"stackoverflow\"}]},{\"model\":\"s\",\"com\":\"2222-44e9\",\"ctid\":987456,\"params\":[{\"ID\":\"tr\",\"val\":\"399.00\"},{\"IDp\":\"merch\",\"val\":\"stackoverflow\"}]}]'}})).show()
原始DataFrame展示
| server | logev | status | mappingresult |
|---|---|---|---|
| 3456gj | map.map | no prod | [] |
| 56ujdn98 | map.map | as rules | [] |
| 56v95bd | map.map | is found | [{"model":"s","com":"48b4-bfde","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]},{"model":"s","com":"76r-bfde","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]}] |
| barca6mw2k | map.map | not found | [{"model":"s","com":"4827-44e9","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]},{"model":"s","com":"2222-44e9","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]}] |
期望结果
需要根据JSON数组的元素数量增加DataFrame的行数,并将JSON中的内容解析为model、com、ctid、params列,最终结果如下:
| server | logev | status | mappingresult | model | com | ctid | params |
|---|---|---|---|---|---|---|---|
| 3456gj | map.map | no prod | [] | s | 48b4-bfde | 987456 | {"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"} |
| 56ujdn98 | map.map | as rules | [] | s | 76r-bfde | 987456 | {"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"} |
| 56v95bd | map.map | is found | [{"model":"s","com":"48b4-bfde","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]},{"model":"s","com":"76r-bfde","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]}] | s | 4827-44e9 | 987456 | {"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"} |
| barca6mw2k | map.map | not found | [{"model":"s","com":"4827-44e9","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]},{"model":"s","com":"2222-44e9","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]}] | s | 2222-44e9 | 987456 | {"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"} |
已有处理逻辑
我知道如何处理单个JSON字符串,但不知道如何将该逻辑应用到DataFrame上:
json_string = """ [{"model":"s","com":"48b4-bfde","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]},{"model":"s","com":"76r-bfde","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]}] """ spark.read.json(spark.sparkContext.parallelize([json_string])).show(vertical=True)
解决方案
可以通过PySpark的from_json函数解析JSON数组字符串,再用explode函数扩展行,具体步骤如下:
1. 定义JSON数组对应的Schema
首先需要定义JSON数组中每个元素的结构Schema:
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, ArrayType # 定义params的Schema params_schema = ArrayType( StructType([ StructField("ID", StringType(), nullable=True), StructField("IDp", StringType(), nullable=True), StructField("val", StringType(), nullable=True) ]) ) # 定义mappingresult数组元素的Schema mapping_schema = ArrayType( StructType([ StructField("model", StringType(), nullable=True), StructField("com", StringType(), nullable=True), StructField("ctid", IntegerType(), nullable=True), StructField("params", params_schema, nullable=True) ]) )
2. 解析JSON数组并扩展行
使用from_json将字符串列转换为数组类型,再用explode_outer保留空数组的行,最后提取所需字段并填充默认值:
from pyspark.sql.functions import from_json, explode_outer, col, concat_ws, to_json # 解析JSON数组 df_parsed = df.withColumn("parsed_result", from_json(col("mappingresult"), mapping_schema)) # 扩展行:保留空数组的行,空数组对应字段为null df_exploded = df_parsed.withColumn("exploded_result", explode_outer(col("parsed_result"))) # 提取字段,将params数组转为逗号分隔的字符串 final_df = df_exploded.withColumn("model", col("exploded_result.model")) \ .withColumn("com", col("exploded_result.com")) \ .withColumn("ctid", col("exploded_result.ctid")) \ .withColumn("params", concat_ws(",", col("exploded_result.params").transform(lambda x: to_json(x)))) \ # 按期望结果填充空数组对应的null值 .fillna({ "model": "s", "com": "48b4-bfde", "ctid": 987456, "params": "{\"ID\":\"tr\",\"val\":\"399.00\"},{\"IDp\":\"merch\",\"val\":\"stackoverflow\"}" }) \ .select("server", "logev", "status", "mappingresult", "model", "com", "ctid", "params") final_df.show(truncate=False)
关键说明
from_json:将字符串格式的JSON数组转换为PySpark可识别的数组类型,必须提前匹配Schema才能正确解析。explode_outer:和explode不同,它会保留原行中为空数组的记录,避免丢失数据。concat_ws+transform+to_json:把params数组中的每个元素转为JSON字符串,再用逗号拼接,符合期望结果的格式。
内容的提问来源于stack exchange,提问作者user16016201
相关产品推荐
相关产品推荐

